huzzle-labs/visual_memory
0
1#!/usr/bin/env python32"""3Evaluation Runner — run an LLM agent against Visual Memory gym scenarios.4 5Single-gym version of the repo-level run_eval.py, tailored for the6visual_memory environment. No --gym flag needed.7 8Usage:9 # Single model (backward compatible)10 python run_eval.py --model gpt-5.4 --save --trajectory11 12 # Multiple models in parallel13 python run_eval.py --model gpt-5.4,claude-sonnet-4-6 --parallel-models 3 --save --trajectory14 15 # Specific scenario16 python run_eval.py --model gpt-5.4 --scenario directional_trap_8x817 18 # pass@k evaluation (run each scenario 10 times, report pass@1, pass@3, pass@8)19 python run_eval.py --model gpt-5.4 --num-samples 10 --pass-k 1,3,8 --save20 21 # Parallel scenarios (run 4 scenarios concurrently per model)22 python run_eval.py --model gpt-5.4 --parallel-scenarios 4 --save23 24 # Resume interrupted run25 python run_eval.py --model gpt-5.4 --run-id my_run --resume --save --trajectory26 27 # ATIF trajectory format (Harbor/Terminus-2 standard)28 python run_eval.py --model gpt-5.4 --trajectory --trajectory-format atif29 30Prerequisites:31 1. pip install -e .32 2. docker build -t openenv-visual-memory -f server/Dockerfile .33 3. docker run -d --name visual-memory -p 8000:8000 openenv-visual-memory34"""35 36import argparse37import json38import logging39import os40import sys41import threading42import time43from concurrent.futures import ThreadPoolExecutor, as_completed44from datetime import datetime, timezone, timedelta45from typing import Any, Dict, List, Optional, Set, Tuple46 47import numpy as np48 49IST = timezone(timedelta(hours=5, minutes=30))50 51from dotenv import load_dotenv52 53load_dotenv(os.path.join(os.path.dirname(__file__), ".env"))54 55sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))56 57from openenv import AutoEnv58 59from agent.runner import AgentRunner60from rewards.base import RewardBreakdown61from rewards.checks import VisualMemoryChecker62from rewards.transforms import VisualMemoryStepTransform63from scenarios.definitions import VISUAL_MEMORY_SCENARIOS64 65logger = logging.getLogger(__name__)66 67GYM_NAME = "visual_memory"68OUTPUT_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "outputs")69 70 71# ── Helpers ──72 73 74def _resolve_base_url() -> str:75 import importlib.resources76 import yaml77 78 try:79 ref = importlib.resources.files(GYM_NAME).joinpath("openenv.yaml")80 with importlib.resources.as_file(ref) as f:81 manifest = yaml.safe_load(f.read_text())82 port = manifest.get("port", 8000)83 return f"http://localhost:{port}"84 except Exception:85 logger.warning("Could not read openenv.yaml, defaulting to port 8000")86 return "http://localhost:8000"87 88 89def _fetch_gym_metadata(base_url: str) -> dict | None:90 import httpx91 92 try:93 resp = httpx.get(f"{base_url}/metadata", timeout=5.0)94 resp.raise_for_status()95 data = resp.json()96 data.pop("readme_content", None)97 return data98 except Exception as e:99 logger.debug(f"Failed to fetch /metadata from {base_url}: {e}")100 return None101 102 103def divider(text: str = ""):104 print(f"\n{'=' * 70}")105 if text:106 print(f" {text}")107 print(f"{'=' * 70}")108 109 110def print_breakdown(breakdown: RewardBreakdown):111 print(breakdown.summary())112 print()113 print(f" Details: {breakdown.details}")114 115 116def _check_label(check: dict) -> str:117 for key in ("min_score", "min_pct", "max_hits"):118 if key in check and key != "type":119 return str(check[key])120 return check.get("type", "?")121 122 123def _short_json(obj, max_len=80):124 s = json.dumps(obj, default=str)125 return s if len(s) <= max_len else s[:max_len] + "..."126 127 128# ── pass@k Estimator ──129 130 131def pass_at_k_estimator(n: int, c: int, k: int) -> float:132 """133 Unbiased estimator of pass@k (Chen et al., 2021 — HumanEval).134 n = total samples, c = correct samples, k = subset size.135 Returns P(at least 1 correct in a random k-subset of n runs).136 """137 if n - c < k:138 return 1.0139 return 1.0 - np.prod(1.0 - k / np.arange(n - c + 1, n + 1))140 141 142# ── Checkpoint (Resume) ──143 144 145def _load_checkpoint(run_id: str, model: str) -> Tuple[Set[str], List[Dict]]:146 """Load completed scenario IDs and their results from a checkpoint file."""147 safe_model = model.replace("/", "_").replace(":", "_")148 checkpoint_path = os.path.join(149 OUTPUT_DIR, "trajectories", run_id, f"{safe_model}.checkpoint.json"150 )151 final_path = os.path.join(152 OUTPUT_DIR, "trajectories", run_id, f"{safe_model}.json"153 )154 155 for path in [checkpoint_path, final_path]:156 if not os.path.exists(path):157 continue158 try:159 with open(path) as f:160 data = json.load(f)161 completed = set()162 prior_results = []163 for s in data.get("scenarios", []):164 sid = s.get("scenario_id")165 reward = s.get("reward")166 if sid and reward is not None:167 completed.add(sid)168 prior_results.append({169 "scenario": sid,170 "total_reward": reward["total"],171 "breakdown": RewardBreakdown(172 structural=reward["structural"],173 ground_truth=reward["ground_truth"],174 efficiency=reward["efficiency"],175 penalty=reward["penalty"],176 total=reward["total"],177 ),178 "steps": s.get("total_steps", 0),179 "elapsed": s.get("elapsed_s", 0),180 "from_checkpoint": True,181 })182 print(f" [{model}] Checkpoint loaded: {len(completed)} scenarios completed")183 return completed, prior_results184 except (json.JSONDecodeError, KeyError) as e:185 logger.warning(f"Could not parse checkpoint {path}: {e}")186 187 return set(), []188 189 190def _build_scenario_entry(r: Dict, scenario) -> Dict:191 """Build a trajectory scenario entry from a result dict."""192 entry = {193 "scenario_id": r.get("scenario", getattr(scenario, "id", "unknown")),194 "elapsed_s": round(r.get("elapsed", 0), 2),195 }196 197 if scenario:198 entry["prompt"] = scenario.prompt199 entry["expected_tools"] = scenario.expected_tools200 entry["max_steps"] = scenario.max_steps201 202 episode = r.get("episode")203 if episode:204 steps = []205 for i, step in enumerate(episode.steps, 1):206 result_data = step.result207 if isinstance(result_data, str):208 try:209 result_data = json.loads(result_data)210 except (json.JSONDecodeError, TypeError):211 pass212 steps.append({213 "step": i,214 "timestamp": step.timestamp,215 "tool_name": step.tool_name,216 "arguments": step.arguments,217 "success": step.success,218 "result": result_data,219 "error": step.error,220 "elapsed_s": round(step.elapsed, 3),221 })222 entry["steps"] = steps223 entry["total_steps"] = len(steps)224 else:225 entry["steps"] = []226 entry["total_steps"] = r.get("steps", 0)227 if r.get("error"):228 entry["error"] = r["error"]229 230 outcome_results = r.get("outcome_results", [])231 if outcome_results and scenario:232 checks = []233 for check_def, passed in zip(scenario.outcome_checks, outcome_results):234 checks.append({"check": check_def, "passed": passed})235 entry["outcome_checks"] = checks236 237 bd = r.get("breakdown")238 if bd:239 entry["reward"] = {240 "structural": round(bd.structural, 4),241 "ground_truth": round(bd.ground_truth, 4),242 "efficiency": round(bd.efficiency, 4),243 "penalty": round(bd.penalty, 4),244 "total": round(bd.total, 4),245 }246 else:247 entry["reward"] = None248 249 # pass@k multi-sample data250 if "samples" in r:251 entry["num_samples"] = r.get("n", 1)252 entry["correct_count"] = r.get("c", 0)253 entry["pass_at_k"] = r.get("pass_at_k", {})254 entry["samples"] = []255 for sample in r["samples"]:256 sample_entry = {257 "sample_idx": sample.get("sample_idx", 0),258 "success": sample.get("success", False),259 "total_reward": sample.get("total_reward", 0),260 "steps": sample.get("steps", 0),261 "elapsed": round(sample.get("elapsed", 0), 2),262 }263 if sample.get("error"):264 sample_entry["error"] = sample["error"]265 entry["samples"].append(sample_entry)266 267 return entry268 269 270def _save_checkpoint(271 run_id: str,272 model: str,273 all_results: List[Dict],274 scenarios: list,275 temperature: float,276 reward_mode: str,277 gym_version: str,278):279 """Incrementally save checkpoint after each scenario."""280 safe_model = model.replace("/", "_").replace(":", "_")281 traj_dir = os.path.join(OUTPUT_DIR, "trajectories", run_id)282 os.makedirs(traj_dir, exist_ok=True)283 checkpoint_path = os.path.join(traj_dir, f"{safe_model}.checkpoint.json")284 285 scenario_map = {s.id: s for s in scenarios}286 checkpoint = {287 "run_id": run_id,288 "model": model,289 "gym": GYM_NAME,290 "gym_version": gym_version,291 "timestamp": datetime.now(IST).isoformat(),292 "temperature": temperature,293 "reward_mode": reward_mode,294 "total_scenarios": len(all_results),295 "scenarios": [],296 }297 298 for r in all_results:299 sid = r.get("scenario")300 if r.get("from_checkpoint"):301 checkpoint["scenarios"].append({302 "scenario_id": sid,303 "elapsed_s": r.get("elapsed", 0),304 "total_steps": r.get("steps", 0),305 "reward": {306 "structural": r["breakdown"].structural,307 "ground_truth": r["breakdown"].ground_truth,308 "efficiency": r["breakdown"].efficiency,309 "penalty": r["breakdown"].penalty,310 "total": r["breakdown"].total,311 } if r.get("breakdown") else None,312 })313 else:314 entry = _build_scenario_entry(r, scenario_map.get(sid))315 checkpoint["scenarios"].append(entry)316 317 with open(checkpoint_path, "w") as f:318 json.dump(checkpoint, f, indent=2, default=str)319 320 321# ── Results & Trajectory Saving ──322 323 324def save_results_to_markdown(325 results: List[Dict[str, Any]],326 model: str,327 output_path: str,328 total_elapsed: float,329 temperature: float,330 run_id: str = "",331 reward_mode: str = "custom",332 gym_version: str = "unknown",333 num_samples: int = 1,334 pass_k_values: Optional[List[int]] = None,335):336 os.makedirs(os.path.dirname(output_path), exist_ok=True)337 338 timestamp = datetime.now(IST).strftime("%Y-%m-%d %H:%M:%S")339 is_new_file = not os.path.exists(output_path)340 is_passk = num_samples > 1 and pass_k_values341 342 with open(output_path, "a") as f:343 if is_new_file:344 f.write(f"# Visual Memory Gym — Evaluation Results\n\n")345 f.write(f"**Run ID**: `{run_id}` \n")346 f.write(f"**Gym Version**: `{gym_version}`\n\n")347 f.write(f"Evaluation results for the **visual_memory** gym across different LLM models.\n\n")348 if is_passk:349 f.write(f"**Mode**: pass@k (n={num_samples}, k={pass_k_values})\n\n")350 if reward_mode == "openenv":351 f.write(f"**Reward Mode**: `openenv` — per-step rewards from `rewards/transforms.py` + ground truth\n\n")352 else:353 f.write(f"**Reward Mode**: `custom` — episode-level rewards from `rewards/base.py`\n\n")354 f.write(f"Trajectories: `outputs/trajectories/{run_id}/`\n\n")355 f.write(f"---\n\n")356 357 safe_model = model.replace("/", "_").replace(":", "_")358 f.write(f"## Model: `{model}`\n\n")359 f.write(f"- **Date**: {timestamp}\n")360 f.write(f"- **Temperature**: {temperature}\n")361 f.write(f"- **Reward Mode**: {reward_mode}\n")362 if is_passk:363 f.write(f"- **Samples per scenario**: {num_samples}\n")364 f.write(f"- **Total Time**: {total_elapsed:.1f}s\n")365 f.write(f"- **Trajectory**: `outputs/trajectories/{run_id}/{safe_model}.json`\n\n")366 367 if is_passk:368 k_headers = " | ".join(f"pass@{k}" for k in pass_k_values)369 f.write(f"| Scenario | n | c | {k_headers} | Best Reward | Avg Steps |\n")370 k_divs = " | ".join(":---:" for _ in pass_k_values)371 f.write(f"|---|:---:|:---:|{k_divs}|:---:|:---:|\n")372 373 for r in results:374 n = r.get("n", 1)375 c = r.get("c", 0)376 pass_at = r.get("pass_at_k", {})377 k_vals = " | ".join(f"{pass_at.get(str(k), 0.0):.2f}" for k in pass_k_values)378 best = r.get("best_reward", 0.0)379 avg_steps = r.get("avg_steps", 0)380 f.write(381 f"| {r['scenario']} "382 f"| {n} | {c} | {k_vals} "383 f"| {best:.2f} | {avg_steps:.1f} |\n"384 )385 386 all_pass1 = [r.get("pass_at_k", {}).get("1", 0.0) for r in results]387 avg_pass1 = sum(all_pass1) / len(all_pass1) if all_pass1 else 0.0388 f.write(f"\n**Average pass@1: {avg_pass1:.2f}**\n\n")389 else:390 if reward_mode == "openenv":391 f.write(f"| Scenario | Quality | Ground Truth | Penalty | **Total** | Steps | Time |\n")392 f.write(f"|---|:---:|:---:|:---:|:---:|:---:|:---:|\n")393 else:394 f.write(f"| Scenario | Structural | Ground Truth | Efficiency | Penalty | **Total** | Steps | Time |\n")395 f.write(f"|---|:---:|:---:|:---:|:---:|:---:|:---:|:---:|\n")396 397 total_reward = 0.0398 for r in results:399 bd = r.get("breakdown")400 if bd:401 if reward_mode == "openenv":402 f.write(403 f"| {r['scenario']} "404 f"| {bd.structural:.2f} "405 f"| {bd.ground_truth:.2f} "406 f"| {bd.penalty:.2f} "407 f"| **{bd.total:.2f}** "408 f"| {r['steps']} "409 f"| {r['elapsed']:.1f}s |\n"410 )411 else:412 f.write(413 f"| {r['scenario']} "414 f"| {bd.structural:.2f} "415 f"| {bd.ground_truth:.2f} "416 f"| {bd.efficiency:.2f} "417 f"| {bd.penalty:.2f} "418 f"| **{bd.total:.2f}** "419 f"| {r['steps']} "420 f"| {r['elapsed']:.1f}s |\n"421 )422 total_reward += bd.total423 else:424 cols = "| — | — | — " if reward_mode == "openenv" else "| — | — | — | — "425 f.write(426 f"| {r['scenario']} "427 f"{cols}"428 f"| **ERROR** "429 f"| {r['steps']} "430 f"| {r['elapsed']:.1f}s |\n"431 )432 433 avg = total_reward / len(results) if results else 0.0434 f.write(f"\n**Average Reward: {avg:.2f}**\n\n")435 436 f.write(f"---\n\n")437 438 logger.info(f"Results saved to {output_path}")439 440 441def save_trajectory(442 results: List[Dict[str, Any]],443 scenarios: list,444 model: str,445 temperature: float,446 total_elapsed: float,447 run_id: str = "",448 reward_mode: str = "custom",449 gym_version: str = "unknown",450 num_samples: int = 1,451 pass_k_values: Optional[List[int]] = None,452):453 run_ts = datetime.now(IST).isoformat()454 455 safe_model = model.replace("/", "_").replace(":", "_")456 filename = f"{safe_model}.json"457 458 traj_dir = os.path.join(OUTPUT_DIR, "trajectories", run_id)459 os.makedirs(traj_dir, exist_ok=True)460 filepath = os.path.join(traj_dir, filename)461 462 scenario_map = {s.id: s for s in scenarios}463 464 trajectory = {465 "run_id": run_id or "untagged",466 "model": model,467 "gym": GYM_NAME,468 "gym_version": gym_version,469 "timestamp": run_ts,470 "temperature": temperature,471 "reward_mode": reward_mode,472 "total_elapsed_s": round(total_elapsed, 2),473 "total_scenarios": len(results),474 }475 476 if num_samples > 1:477 trajectory["num_samples"] = num_samples478 trajectory["pass_k_values"] = pass_k_values or [1]479 480 trajectory["scenarios"] = []481 482 for r in results:483 sid = r.get("scenario")484 scenario = scenario_map.get(sid)485 if r.get("from_checkpoint"):486 trajectory["scenarios"].append({487 "scenario_id": sid,488 "elapsed_s": r.get("elapsed", 0),489 "total_steps": r.get("steps", 0),490 "reward": {491 "structural": r["breakdown"].structural,492 "ground_truth": r["breakdown"].ground_truth,493 "efficiency": r["breakdown"].efficiency,494 "penalty": r["breakdown"].penalty,495 "total": r["breakdown"].total,496 } if r.get("breakdown") else None,497 "from_checkpoint": True,498 })499 else:500 entry = _build_scenario_entry(r, scenario)501 trajectory["scenarios"].append(entry)502 503 totals = [504 s["reward"]["total"]505 for s in trajectory["scenarios"]506 if s.get("reward")507 ]508 trajectory["avg_reward"] = round(sum(totals) / len(totals), 4) if totals else 0.0509 510 with open(filepath, "w") as f:511 json.dump(trajectory, f, indent=2, default=str)512 513 print(f"\n Trajectory saved: {filepath}")514 logger.info(f"Trajectory saved to {filepath}")515 516 # Clean up checkpoint file now that full trajectory is written517 checkpoint_path = os.path.join(traj_dir, f"{safe_model}.checkpoint.json")518 if os.path.exists(checkpoint_path):519 os.remove(checkpoint_path)520 521 return filepath522 523 524def save_trajectory_atif(525 results: List[Dict[str, Any]],526 scenarios: list,527 model: str,528 temperature: float,529 total_elapsed: float,530 run_id: str = "",531 reward_mode: str = "custom",532 gym_version: str = "unknown",533 token_usage: Optional[Dict[str, int]] = None,534):535 """Save trajectory in ATIF v1.4 format (Harbor/Terminus-2 standard)."""536 safe_model = model.replace("/", "_").replace(":", "_")537 filename = f"{safe_model}_atif.json"538 traj_dir = os.path.join(OUTPUT_DIR, "trajectories", run_id)539 os.makedirs(traj_dir, exist_ok=True)540 filepath = os.path.join(traj_dir, filename)541 542 scenario_map = {s.id: s for s in scenarios}543 atif_steps = []544 step_id = 0545 546 for r in results:547 if r.get("from_checkpoint"):548 continue549 550 sid = r.get("scenario")551 scenario = scenario_map.get(sid)552 episode = r.get("episode")553 if not episode:554 continue555 556 step_id += 1557 atif_steps.append({558 "step_id": step_id,559 "timestamp": episode.steps[0].timestamp if episode.steps else None,560 "role": "user",561 "message": scenario.prompt if scenario else sid,562 "tool_calls": [],563 "observation": None,564 })565 566 for step in episode.steps:567 step_id += 1568 result_data = step.result569 if isinstance(result_data, str):570 try:571 result_data = json.loads(result_data)572 except (json.JSONDecodeError, TypeError):573 pass574 575 atif_steps.append({576 "step_id": step_id,577 "timestamp": step.timestamp,578 "role": "assistant",579 "message": None,580 "tool_calls": [{581 "id": f"call_{step_id}",582 "function_name": step.tool_name,583 "arguments": step.arguments,584 }],585 "observation": {586 "content": result_data,587 "success": step.success,588 "error": step.error,589 },590 "duration_ms": round(step.elapsed * 1000),591 })592 593 bd = r.get("breakdown")594 if bd and atif_steps:595 atif_steps[-1]["reward"] = bd.total596 597 usage = token_usage or {}598 trajectory = {599 "schema_version": "ATIF-v1.4",600 "session_id": run_id,601 "started_at": datetime.now(IST).isoformat(),602 "agent": {603 "name": f"openenv-{GYM_NAME}",604 "model_name": model,605 "temperature": temperature,606 },607 "environment": {608 "gym": GYM_NAME,609 "version": gym_version,610 "reward_mode": reward_mode,611 },612 "steps": atif_steps,613 "final_metrics": {614 "total_steps": step_id,615 "total_wall_time_s": round(total_elapsed, 2),616 "total_prompt_tokens": usage.get("prompt_tokens", 0),617 "total_completion_tokens": usage.get("completion_tokens", 0),618 "total_cost_usd": 0.0,619 "custom": {620 "avg_reward": round(621 sum(r.get("total_reward", 0) for r in results) / max(len(results), 1), 4622 ),623 "total_scenarios": len(results),624 },625 },626 }627 628 with open(filepath, "w") as f:629 json.dump(trajectory, f, indent=2, default=str)630 631 print(f"\n ATIF trajectory saved: {filepath}")632 logger.info(f"ATIF trajectory saved to {filepath}")633 return filepath634 635 636# ── Scenario Execution ──637 638 639WS_RETRY_ERRORS = ("ConnectionClosed", "ConnectionClosedOK", "ConnectionClosedError", "sent 1000")640MAX_WS_RETRIES = 3641 642 643def _run_scenario_with_retries(644 scenario,645 runner: AgentRunner,646 checker,647 env_client,648 connect_fn,649 model: str,650) -> Dict[str, Any]:651 """Execute a single scenario with WebSocket retry logic. Returns a result dict."""652 start = time.time()653 last_error = None654 655 for attempt in range(MAX_WS_RETRIES + 1):656 try:657 if attempt > 0:658 logger.info(f"[{model}] Reconnecting (attempt {attempt + 1}) for {scenario.id}")659 print(f" [{model}] Reconnecting WebSocket (attempt {attempt + 1})...")660 try:661 env_client.__exit__(None, None, None)662 except Exception:663 pass664 time.sleep(2 * attempt)665 env_client, runner = connect_fn()666 667 episode, breakdown = runner.run_scenario(scenario, checker)668 elapsed = time.time() - start669 670 if hasattr(checker, "set_episode"):671 checker.set_episode(episode)672 outcome_results = checker.check_all(scenario.outcome_checks)673 674 return {675 "scenario": scenario.id,676 "total_reward": breakdown.total,677 "breakdown": breakdown,678 "steps": len(episode.steps),679 "elapsed": elapsed,680 "episode": episode,681 "outcome_results": outcome_results,682 }683 684 except Exception as e:685 last_error = e686 is_ws_error = any(tok in type(e).__name__ or tok in str(e) for tok in WS_RETRY_ERRORS)687 if is_ws_error and attempt < MAX_WS_RETRIES:688 logger.warning(f"[{model}] WebSocket error on {scenario.id}: {e}")689 continue690 break691 692 elapsed = time.time() - start693 logger.exception(f"[{model}] Scenario {scenario.id} failed")694 return {695 "scenario": scenario.id,696 "total_reward": 0.0,697 "breakdown": None,698 "steps": 0,699 "elapsed": elapsed,700 "error": str(last_error),701 }702 703 704def _run_scenario_n_samples(705 scenario,706 n: int,707 pass_threshold: float,708 pass_k_values: List[int],709 model: str,710 base_url: str,711 temperature: float,712 max_tokens: int,713 reward_mode: str,714) -> Dict[str, Any]:715 """Run a single scenario n times for pass@k evaluation."""716 samples = []717 correct_count = 0718 719 def _connect():720 client = AutoEnv.from_env(GYM_NAME, base_url=base_url)721 client.__enter__()722 xform = VisualMemoryStepTransform() if reward_mode == "openenv" else None723 rnr = AgentRunner(724 model=model,725 env_client=client,726 temperature=temperature,727 max_tokens=max_tokens,728 reward_mode=reward_mode,729 transform=xform,730 )731 return client, rnr732 733 env_client, runner = _connect()734 checker = VisualMemoryChecker()735 736 try:737 for sample_idx in range(n):738 result = _run_scenario_with_retries(739 scenario, runner, checker, env_client, _connect, model,740 )741 742 gt_score = 0.0743 outcome_results = result.get("outcome_results", [])744 if outcome_results:745 gt_score = sum(outcome_results) / len(outcome_results)746 is_success = gt_score >= pass_threshold747 748 if is_success:749 correct_count += 1750 751 result["sample_idx"] = sample_idx752 result["ground_truth_score"] = gt_score753 result["success"] = is_success754 samples.append(result)755 756 status = "PASS" if is_success else "FAIL"757 print(758 f" [{model}] {scenario.id} sample {sample_idx + 1}/{n}: "759 f"{status} (gt={gt_score:.2f}, reward={result['total_reward']:.2f}, "760 f"{result['steps']} steps, {result['elapsed']:.1f}s)"761 )762 finally:763 try:764 env_client.__exit__(None, None, None)765 except Exception:766 pass767 768 pass_at_k = {}769 for k in pass_k_values:770 if k <= n:771 pass_at_k[str(k)] = round(pass_at_k_estimator(n, correct_count, k), 4)772 773 best_sample = max(samples, key=lambda s: s.get("total_reward", 0.0))774 775 return {776 "scenario": scenario.id,777 "n": n,778 "c": correct_count,779 "pass_at_k": pass_at_k,780 "samples": samples,781 "best_reward": best_sample.get("total_reward", 0.0),782 "avg_steps": sum(s.get("steps", 0) for s in samples) / max(len(samples), 1),783 "total_reward": best_sample.get("total_reward", 0.0),784 "breakdown": best_sample.get("breakdown"),785 "steps": best_sample.get("steps", 0),786 "elapsed": sum(s.get("elapsed", 0) for s in samples),787 "episode": best_sample.get("episode"),788 "outcome_results": best_sample.get("outcome_results", []),789 }790 791 792# ── Model Workers ──793 794 795def _run_single_model(796 model: str,797 base_url: str,798 scenarios: list,799 temperature: float,800 max_tokens: int,801 reward_mode: str,802 run_id: str,803 save: bool,804 trajectory: bool,805 verbose: bool,806 gym_version: str = "unknown",807 num_samples: int = 1,808 pass_k_values: Optional[List[int]] = None,809 pass_threshold: float = 0.5,810 parallel_scenarios: int = 1,811 resume: bool = False,812 trajectory_format: str = "native",813) -> Dict[str, Any]:814 model_start = time.time()815 816 # Resume: load checkpoint817 completed_ids: Set[str] = set()818 prior_results: List[Dict] = []819 if resume:820 completed_ids, prior_results = _load_checkpoint(run_id, model)821 822 pending = [s for s in scenarios if s.id not in completed_ids]823 if not pending:824 print(f" [{model}] All scenarios already completed (checkpoint)")825 model_elapsed = time.time() - model_start826 return {"model": model, "results": prior_results, "elapsed": model_elapsed}827 828 if completed_ids:829 print(f" [{model}] Resuming: {len(pending)} remaining of {len(scenarios)} scenarios")830 831 model_results = list(prior_results)832 results_lock = threading.Lock()833 834 is_passk = num_samples > 1835 836 def _connect():837 client = AutoEnv.from_env(GYM_NAME, base_url=base_url)838 client.__enter__()839 xform = VisualMemoryStepTransform() if reward_mode == "openenv" else None840 rnr = AgentRunner(841 model=model,842 env_client=client,843 temperature=temperature,844 max_tokens=max_tokens,845 reward_mode=reward_mode,846 transform=xform,847 )848 return client, rnr849 850 def _run_one_scenario(scenario, idx, total):851 """Run a single scenario (optionally n samples) and append to results."""852 print(f"\n [{model}] Scenario {idx}/{total}: {scenario.id}")853 854 if is_passk:855 result = _run_scenario_n_samples(856 scenario, n=num_samples, pass_threshold=pass_threshold,857 pass_k_values=pass_k_values or [1], model=model,858 base_url=base_url, temperature=temperature,859 max_tokens=max_tokens, reward_mode=reward_mode,860 )861 pk = result.get("pass_at_k", {})862 pk_str = ", ".join(f"pass@{k}={pk.get(str(k), 0):.2f}" for k in (pass_k_values or [1]))863 print(864 f" [{model}] {scenario.id}: {result['c']}/{result['n']} correct → {pk_str}"865 )866 else:867 env_client, runner = _connect()868 checker = VisualMemoryChecker()869 try:870 result = _run_scenario_with_retries(871 scenario, runner, checker, env_client, _connect, model,872 )873 reward_str = f"{result['total_reward']:.2f}" if result.get("breakdown") else "ERROR"874 print(875 f" [{model}] {scenario.id}: {reward_str} "876 f"({result['steps']} steps, {result['elapsed']:.1f}s)"877 )878 finally:879 try:880 env_client.__exit__(None, None, None)881 except Exception:882 pass883 884 with results_lock:885 model_results.append(result)886 _save_checkpoint(887 run_id, model, model_results, scenarios,888 temperature, reward_mode, gym_version,889 )890 891 return result892 893 if parallel_scenarios > 1 and len(pending) > 1:894 max_concurrent = int(os.getenv("MAX_CONCURRENT_ENVS", "8"))895 max_workers = min(parallel_scenarios, len(pending), max_concurrent)896 print(f" [{model}] Running {len(pending)} scenarios with {max_workers} parallel workers")897 898 with ThreadPoolExecutor(max_workers=max_workers) as executor:899 futures = {}900 for idx, scenario in enumerate(pending, len(completed_ids) + 1):901 future = executor.submit(902 _run_one_scenario, scenario, idx, len(scenarios),903 )904 futures[future] = scenario905 906 for future in as_completed(futures):907 scenario = futures[future]908 try:909 future.result()910 except Exception as e:911 print(f" [{model}] {scenario.id}: ERROR - {e}")912 logger.exception(f"Scenario {scenario.id} failed")913 with results_lock:914 model_results.append({915 "scenario": scenario.id,916 "total_reward": 0.0,917 "breakdown": None,918 "steps": 0,919 "elapsed": 0.0,920 "error": str(e),921 })922 else:923 if not is_passk:924 env_client, runner = _connect()925 checker = VisualMemoryChecker()926 try:927 for idx, scenario in enumerate(pending, len(completed_ids) + 1):928 print(f"\n [{model}] Scenario {idx}/{len(scenarios)}: {scenario.id}")929 result = _run_scenario_with_retries(930 scenario, runner, checker, env_client, _connect, model,931 )932 reward_str = f"{result['total_reward']:.2f}" if result.get("breakdown") else "ERROR"933 print(934 f" [{model}] {scenario.id}: {reward_str} "935 f"({result['steps']} steps, {result['elapsed']:.1f}s)"936 )937 model_results.append(result)938 _save_checkpoint(939 run_id, model, model_results, scenarios,940 temperature, reward_mode, gym_version,941 )942 finally:943 try:944 env_client.__exit__(None, None, None)945 except Exception:946 pass947 else:948 for idx, scenario in enumerate(pending, len(completed_ids) + 1):949 _run_one_scenario(scenario, idx, len(scenarios))950 951 model_elapsed = time.time() - model_start952 953 if save:954 output_path = os.path.join(OUTPUT_DIR, "results", f"{run_id}.md")955 save_results_to_markdown(956 results=model_results,957 model=model,958 output_path=output_path,959 total_elapsed=model_elapsed,960 temperature=temperature,961 run_id=run_id,962 reward_mode=reward_mode,963 gym_version=gym_version,964 num_samples=num_samples,965 pass_k_values=pass_k_values,966 )967 968 if trajectory:969 save_trajectory(970 results=model_results,971 scenarios=scenarios,972 model=model,973 temperature=temperature,974 total_elapsed=model_elapsed,975 run_id=run_id,976 reward_mode=reward_mode,977 gym_version=gym_version,978 num_samples=num_samples,979 pass_k_values=pass_k_values,980 )981 if trajectory_format == "atif":982 save_trajectory_atif(983 results=model_results,984 scenarios=scenarios,985 model=model,986 temperature=temperature,987 total_elapsed=model_elapsed,988 run_id=run_id,989 reward_mode=reward_mode,990 gym_version=gym_version,991 )992 993 return {994 "model": model,995 "results": model_results,996 "elapsed": model_elapsed,997 }998 999 1000def _run_single_model_detailed(1001 model: str,1002 base_url: str,1003 scenarios: list,1004 temperature: float,1005 max_tokens: int,1006 reward_mode: str,1007 run_id: str,1008 save: bool,1009 trajectory: bool,1010 gym_version: str = "unknown",1011 num_samples: int = 1,1012 pass_k_values: Optional[List[int]] = None,1013 pass_threshold: float = 0.5,1014 resume: bool = False,1015 trajectory_format: str = "native",1016) -> Dict[str, Any]:1017 model_start = time.time()1018 results = []1019 1020 completed_ids: Set[str] = set()1021 if resume:1022 completed_ids, prior = _load_checkpoint(run_id, model)1023 results = list(prior)1024 1025 pending = [s for s in scenarios if s.id not in completed_ids]1026 if not pending:1027 print(f" [{model}] All scenarios already completed (checkpoint)")1028 return {"model": model, "results": results, "elapsed": time.time() - model_start}1029 1030 is_passk = num_samples > 11031 1032 if is_passk:1033 for i, scenario in enumerate(pending, len(completed_ids) + 1):1034 divider(f"Scenario {i}/{len(scenarios)}: {scenario.id}")1035 print(f" Prompt: {scenario.prompt[:120]}...")1036 print(f" Samples: {num_samples}")1037 print()1038 1039 result = _run_scenario_n_samples(1040 scenario, n=num_samples, pass_threshold=pass_threshold,1041 pass_k_values=pass_k_values or [1], model=model,1042 base_url=base_url, temperature=temperature,1043 max_tokens=max_tokens, reward_mode=reward_mode,1044 )1045 1046 pk = result.get("pass_at_k", {})1047 print(f"\n -- pass@k Results --")1048 print(f" Correct: {result['c']}/{result['n']}")1049 for k in (pass_k_values or [1]):1050 print(f" pass@{k}: {pk.get(str(k), 0.0):.4f}")1051 1052 results.append(result)1053 _save_checkpoint(1054 run_id, model, results, scenarios,1055 temperature, reward_mode, gym_version,1056 )1057 else:1058 env_client = AutoEnv.from_env(GYM_NAME, base_url=base_url)1059 env_client.__enter__()1060 checker = VisualMemoryChecker()1061 transform = VisualMemoryStepTransform() if reward_mode == "openenv" else None1062 runner = AgentRunner(1063 model=model, env_client=env_client, temperature=temperature,1064 max_tokens=max_tokens, reward_mode=reward_mode, transform=transform,1065 )1066 1067 try:1068 for i, scenario in enumerate(pending, len(completed_ids) + 1):1069 divider(f"Scenario {i}/{len(scenarios)}: {scenario.id}")1070 print(f" Prompt: {scenario.prompt[:120]}...")1071 print(f" Expected tools: {scenario.expected_tools}")1072 print(f" Max steps: {scenario.max_steps}")1073 print()1074 1075 start = time.time()1076 try:1077 episode, breakdown = runner.run_scenario(scenario, checker)1078 elapsed = time.time() - start1079 1080 print()1081 print(" -- Agent Actions --")1082 for step in episode.steps:1083 status = "OK" if step.success else "FAIL"1084 args_str = _short_json(step.arguments)1085 print(f" [{status}] {step.tool_name}({args_str})")1086 print(f" Steps taken: {len(episode.steps)}")1087 1088 if hasattr(checker, "set_episode"):1089 checker.set_episode(episode)1090 1091 print()1092 print(" -- Ground Truth Verification --")1093 outcome_results = checker.check_all(scenario.outcome_checks)1094 for check, score in zip(scenario.outcome_checks, outcome_results):1095 status = "PASS" if score else "FAIL"1096 label = _check_label(check)1097 print(f" [{status}] {check['type']}: {label}")1098 1099 print()1100 print(" -- Reward Breakdown --")1101 print_breakdown(breakdown)1102 print(f"\n Completed in {elapsed:.1f}s")1103 1104 result = {1105 "scenario": scenario.id,1106 "total_reward": breakdown.total,1107 "breakdown": breakdown,1108 "steps": len(episode.steps),1109 "elapsed": elapsed,1110 "episode": episode,1111 "outcome_results": outcome_results,1112 }1113 results.append(result)1114 1115 except Exception as e:1116 elapsed = time.time() - start1117 print(f"\n ERROR: {e}")1118 logger.exception(f"Scenario {scenario.id} failed")1119 results.append({1120 "scenario": scenario.id,1121 "total_reward": 0.0,1122 "breakdown": None,1123 "steps": 0,1124 "elapsed": elapsed,1125 "error": str(e),1126 })1127 1128 _save_checkpoint(1129 run_id, model, results, scenarios,1130 temperature, reward_mode, gym_version,1131 )1132 1133 finally:1134 env_client.__exit__(None, None, None)1135 logger.info("AutoEnv client disconnected.")1136 1137 model_elapsed = time.time() - model_start1138 1139 if save:1140 output_path = os.path.join(OUTPUT_DIR, "results", f"{run_id}.md")1141 save_results_to_markdown(1142 results=results, model=model, output_path=output_path,1143 total_elapsed=model_elapsed, temperature=temperature,1144 run_id=run_id, reward_mode=reward_mode, gym_version=gym_version,1145 num_samples=num_samples, pass_k_values=pass_k_values,1146 )1147 print(f"\n Results saved: {output_path}")1148 1149 if trajectory:1150 save_trajectory(1151 results=results, scenarios=scenarios, model=model,1152 temperature=temperature, total_elapsed=model_elapsed,1153 run_id=run_id, reward_mode=reward_mode, gym_version=gym_version,1154 num_samples=num_samples, pass_k_values=pass_k_values,1155 )1156 if trajectory_format == "atif":1157 save_trajectory_atif(1158 results=results, scenarios=scenarios, model=model,1159 temperature=temperature, total_elapsed=model_elapsed,1160 run_id=run_id, reward_mode=reward_mode, gym_version=gym_version,1161 )1162 1163 return {1164 "model": model,1165 "results": results,1166 "elapsed": model_elapsed,1167 }1168 1169 1170# ── Main ──1171 1172 1173def main():1174 parser = argparse.ArgumentParser(1175 description="Evaluate an LLM agent against Visual Memory gym scenarios.",1176 formatter_class=argparse.RawDescriptionHelpFormatter,1177 epilog="""1178Examples:1179 # Basic (backward compatible)1180 python run_eval.py --model gpt-5.4 --save --trajectory1181 python run_eval.py --model gpt-5.4 --scenario directional_trap_8x81182 1183 # pass@k1184 python run_eval.py --model gpt-5.4 --num-samples 10 --pass-k 1,3,8 --save1185 1186 # Parallel scenarios1187 python run_eval.py --model gpt-5.4 --parallel-scenarios 4 --save1188 1189 # Resume interrupted run1190 python run_eval.py --model gpt-5.4 --run-id my_run --resume --save --trajectory1191 1192 # ATIF trajectory format1193 python run_eval.py --model gpt-5.4 --trajectory --trajectory-format atif1194 1195 # Combined1196 python run_eval.py --model gpt-5.4 --num-samples 10 --pass-k 1,3,8 \\1197 --parallel-scenarios 4 --run-id bench_v1 --resume --save \\1198 --trajectory --trajectory-format atif1199 """,1200 )