"""BenchFlow 评测执行、完整性判断与恢复。""" from __future__ import annotations import json from pathlib import Path import subprocess import sys from typing import Any import uuid from scripts.dynamic_compile.fast.storage import sha256_file from .manifest import read_json_object from .paths import PROJECT_ROOT EVALUATION_SCRIPT = PROJECT_ROOT / "scripts" / "evaluate" / "run-raw-task.sh" MAX_PARALLEL = 3 def evaluation_artifacts_complete(test_dir: Path) -> bool: """Return whether one BenchFlow attempt has all required usable artifacts.""" summary = read_json_object(test_dir / "summary.json") required_skill = read_json_object(test_dir / "required-skill.json") if summary is None or required_skill is None: return False try: total = int(summary.get("total", 0) or 0) passed = int(summary.get("passed", summary.get("pass", 0)) or 0) failed = int(summary.get("failed", summary.get("fail", 0)) or 0) errored = int(summary.get("errored", summary.get("error", 0)) or 0) verifier_errored = int(summary.get("verifier_errored", 0) or 0) except (TypeError, ValueError): return False if total != 1 or passed + failed != 1 or errored or verifier_errored: return False if required_skill.get("invoked") is not True or required_skill.get("parse_errors"): return False trajectories = sorted(test_dir.rglob("acp_trajectory.jsonl")) canonical = [path for path in trajectories if "trajectory" in path.parts] trajectory = canonical[0] if canonical else (trajectories[0] if trajectories else None) if trajectory is None: return False result = read_json_object(trajectory.parent.parent / "result.json") if result is None: return False try: events = [ json.loads(line) for line in trajectory.read_text(encoding="utf-8").splitlines() if line.strip() ] except (OSError, json.JSONDecodeError): return False return bool(events) and all(isinstance(event, dict) for event in events) def completed_evaluation_count(output: Path) -> int: if not output.is_dir(): return 0 return sum( evaluation_artifacts_complete(test_dir) for test_dir in output.glob("test-*") if test_dir.is_dir() ) def quarantine_incomplete_evaluations(output: Path) -> int: if not output.is_dir(): return 0 incomplete = [ test_dir for test_dir in sorted(output.glob("test-*")) if test_dir.is_dir() and not evaluation_artifacts_complete(test_dir) ] if not incomplete: return 0 quarantine = output / ".incomplete" quarantine.mkdir(exist_ok=True) for test_dir in incomplete: destination = quarantine / test_dir.name if destination.exists(): destination = quarantine / f"{test_dir.name}-{uuid.uuid4().hex[:6]}" test_dir.replace(destination) print( f"[compile-pipeline] preserved incomplete evaluation at {destination}", file=sys.stderr, flush=True, ) return len(incomplete) def evaluate( *, harness: str, model: str, task_dir: Path, skill_source: Path, output: Path, repeat: int, ) -> None: completed = completed_evaluation_count(output) if completed >= repeat: print( f"[compile-pipeline] evaluation complete: {output} ({completed}/{repeat}); skipping", file=sys.stderr, flush=True, ) return quarantine_incomplete_evaluations(output) missing = repeat - completed if completed: print( f"[compile-pipeline] evaluation incomplete: {output} " f"({completed}/{repeat}); running {missing} missing rollout(s)", file=sys.stderr, flush=True, ) command = [ "bash", str(EVALUATION_SCRIPT), "--harness", harness, "--model", model, "--task", str(task_dir), "--skill-source", str(skill_source), "--output", str(output), "--require-skill", "--repeat", str(missing), "--max-parallel", str(MAX_PARALLEL), ] try: subprocess.run(command, cwd=PROJECT_ROOT, check=True) except subprocess.CalledProcessError as exc: raise RuntimeError( f"evaluation failed for {skill_source} with exit code {exc.returncode}" ) from exc completed = completed_evaluation_count(output) if completed < repeat: raise RuntimeError( f"evaluation produced only {completed}/{repeat} complete rollout artifacts " f"under {output}" ) def complete_deep_skill(output: Path) -> Path | None: state = read_json_object(output / "run.json") report = read_json_object(output / "report.json") skill = output / "S_final" if ( state is not None and state.get("status") == "complete" and report is not None and report.get("status") == "complete" and (skill / "SKILL.md").is_file() ): return skill return None def deep_final_rollouts(output: Path, final_skill: Path, repeat: int) -> Path | None: report = read_json_object(output / "report.json") if report is None: return None value = report.get("final_rollouts") expected_hash = report.get("final_rollout_skill_sha256") if not isinstance(value, str) or not isinstance(expected_hash, str): return None rollouts = Path(value).resolve() if ( not rollouts.is_dir() or expected_hash != sha256_file(final_skill / "SKILL.md") or completed_evaluation_count(rollouts) < repeat ): return None return rollouts