dora7's picture
Albedo SN97 workspace v16: RECORD, merged v16, adapters v13/v15/v16-dpo, packs, harness, evals
2abcc30 verified
Raw
History Blame Contribute Delete
10.1 kB
from __future__ import annotations
import asyncio
import json
import os
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from uuid import uuid4
from albedo_config import JudgeSettings
from albedo_config.models import JUDGE_MODELS
from albedo_eval_service.judge_api import (
JudgeSample,
ObservationSimulationService,
QuestionPrepStore,
QuestionService,
ReferenceTrajectoryService,
RepoContextClient,
ScoreBatchRequest,
_score_samples,
)
from albedo_eval_service.judge_core import aggregate_scores, challenger_beats_king
from albedo_eval_service.judge_llm_client import JudgeLLMClient
from .constants import DEFAULT_DATA_ROOT, DEFAULT_RUNS_DIR
from .samples import load_samples
DEFAULT_JUDGE_RUN = Path("/workspace/data/eval-runs/20260824T182728Z-cdd1388b")
SOTA_TURNS = 8
# 8 SOTA + ~7 sims + 4 checklist/behavior + 1 prune + 2 sides
ESTIMATED_CALLS_PER_SAMPLE = 22
def promote_openrouter_key() -> str:
"""Accept OPENROUTER_API_KEY as a synonym for the official judge setting."""
key = (os.environ.get("ALBEDO_JUDGE_OPENROUTER_API_KEY") or "").strip()
if key:
return key
fallback = (os.environ.get("OPENROUTER_API_KEY") or "").strip()
if fallback:
os.environ["ALBEDO_JUDGE_OPENROUTER_API_KEY"] = fallback
return fallback
return ""
def require_openrouter_key() -> str:
key = promote_openrouter_key()
if key:
return key
raise SystemExit(
"Live judge needs ALBEDO_JUDGE_OPENROUTER_API_KEY (OpenRouter, model z-ai/glm-5.2).\n"
"export ALBEDO_JUDGE_OPENROUTER_API_KEY=... && python -m local_eval judge "
f"--run {DEFAULT_JUDGE_RUN} --limit 2"
)
def load_duel_rows(run_dir: Path) -> list[dict[str, Any]]:
path = Path(run_dir) / "generated-samples.jsonl"
if not path.is_file():
raise SystemExit(f"no generated-samples.jsonl under {run_dir}")
rows = []
for line in path.read_text().splitlines():
if line.strip():
rows.append(json.loads(line))
if not rows:
raise SystemExit(f"empty generated-samples.jsonl under {run_dir}")
return rows
def load_verdict(run_dir: Path) -> dict[str, Any]:
path = Path(run_dir) / "verdict.json"
if not path.is_file():
return {}
return json.loads(path.read_text())
def load_judge_samples(
run_dir: Path,
*,
dataset_root: Path = DEFAULT_DATA_ROOT,
seed: str | None = None,
limit: int = 0,
) -> tuple[list[JudgeSample], dict[str, Any], list[str]]:
"""Rebuild official JudgeSamples from a local duel run.
Submit salt is applied on the *full* run sample list so a --limit slice
keeps the same markers the models actually saw.
"""
run_dir = Path(run_dir)
verdict = load_verdict(run_dir)
rows = load_duel_rows(run_dir)
seed = seed or str(verdict.get("seed") or "local-eval")
all_ids = list(verdict.get("sample_ids") or [row["sample_id"] for row in rows])
prefixes = load_samples(
Path(dataset_root),
sample_ids=all_ids,
sample_count=len(all_ids),
seed=seed,
)
by_id = {sample.sample_id: sample for sample in prefixes}
usable = [
row
for row in rows
if row.get("sample_id") in by_id
and row.get("previous_king_output")
and row.get("challenger_output")
and not row.get("king_error")
and not row.get("chal_error")
]
if limit > 0:
usable = usable[:limit]
warnings: list[str] = []
samples: list[JudgeSample] = []
for row in usable:
prefix = by_id[row["sample_id"]]
stored_marker = str(row.get("submit_marker") or "")
if stored_marker and stored_marker != prefix.submit_marker:
warnings.append(
f"{row['sample_id']}: reloaded marker {prefix.submit_marker!r} "
f"!= stored {stored_marker!r}"
)
samples.append(
JudgeSample(
sample_id=prefix.sample_id,
prompt=prefix.prompt,
previous_king_output=row["previous_king_output"],
challenger_output=row["challenger_output"],
messages=prefix.messages,
submit_marker=stored_marker or prefix.submit_marker,
submit_command=str(row.get("submit_command") or prefix.submit_command),
)
)
return samples, verdict, warnings
def preview_samples(samples: list[JudgeSample]) -> list[dict[str, Any]]:
return [
{
"sample_id": sample.sample_id,
"prompt_chars": len(sample.prompt or ""),
"messages": len(sample.messages or []),
"king_chars": len(sample.previous_king_output or ""),
"chal_chars": len(sample.challenger_output or ""),
"submit_marker": sample.submit_marker,
}
for sample in samples
]
def _settings() -> JudgeSettings:
promote_openrouter_key()
return JudgeSettings()
async def _score(request: ScoreBatchRequest, settings: JudgeSettings) -> list[dict[str, Any]]:
client = JudgeLLMClient(settings)
repo = RepoContextClient(settings) if settings.repo_context_url else None
try:
simulator = ObservationSimulationService(settings, client, repo)
questions = QuestionService(
settings,
client,
ReferenceTrajectoryService(settings, client, simulator),
)
store = QuestionPrepStore(settings, questions)
return await _score_samples(
client=client, request=request, settings=settings, prep_store=store
)
finally:
await client.aclose()
if repo is not None:
await repo.aclose()
def run_judge(
*,
run_dir: Path = DEFAULT_JUDGE_RUN,
dataset_root: Path = DEFAULT_DATA_ROOT,
seed: str | None = None,
limit: int = 2,
dry_run: bool = False,
out_dir: Path | None = None,
) -> dict[str, Any]:
samples, verdict, warnings = load_judge_samples(
run_dir, dataset_root=dataset_root, seed=seed, limit=limit
)
if not samples:
raise SystemExit(f"no scorable king+challenger pairs under {run_dir}")
preview = preview_samples(samples)
report: dict[str, Any] = {
"mode": "official_glm_checklist",
"source_run": str(Path(run_dir).resolve()),
"source_run_id": verdict.get("run_id"),
"source_proxy_challenger": verdict.get("score_challenger"),
"source_proxy_king": verdict.get("score_king"),
"seed": seed or verdict.get("seed") or "local-eval",
"judge_models": list(JUDGE_MODELS),
"sota_turns": SOTA_TURNS,
"sample_count": len(samples),
"estimated_llm_calls": len(samples) * ESTIMATED_CALLS_PER_SAMPLE,
"sample_ids": [sample.sample_id for sample in samples],
"preview": preview,
"warnings": warnings,
"note": (
"Official score-batch path: GLM-5.2 SOTA reference (8 turns, LLM-sim observations "
"unless ALBEDO_JUDGE_REPO_CONTEXT_URL is set), checklist, then yes/no judge on "
"both stored trajectories. This is not proxy_score."
),
}
print(json.dumps({k: report[k] for k in report if k != "preview"}, indent=2), flush=True)
for row in preview:
print(
f" {row['sample_id']} prompt={row['prompt_chars']} "
f"king={row['king_chars']} chal={row['chal_chars']} "
f"marker={row['submit_marker']!r}",
flush=True,
)
if dry_run:
report["state"] = "dry_run"
return report
require_openrouter_key()
settings = _settings()
eval_run_id = (
datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + "-judge-" + uuid4().hex[:8]
)
request = ScoreBatchRequest(
eval_run_id=eval_run_id,
batch_id="score-0001",
samples=samples,
total_sample_count=len(samples),
judge_models=list(JUDGE_MODELS),
)
print(
f"live judge start run={eval_run_id} samples={len(samples)} "
f"models={list(JUDGE_MODELS)} (~{report['estimated_llm_calls']} LLM calls)",
flush=True,
)
records = asyncio.run(_score(request, settings))
summary = aggregate_scores(records, min_valid_fraction=settings.min_valid_fraction)
dest = Path(out_dir or run_dir)
dest.mkdir(parents=True, exist_ok=True)
(dest / "official-judge.json").write_text(
json.dumps({**report, "eval_run_id": eval_run_id, **summary}, indent=2) + "\n"
)
with (dest / "official-judge-records.jsonl").open("w") as handle:
for record in records:
handle.write(json.dumps(record, ensure_ascii=False) + "\n")
per_sample = [
{
"sample_id": record.get("sample_id"),
"scored": record.get("scored"),
"king_score": record.get("king_score"),
"challenger_score": record.get("challenger_score"),
"error": record.get("error"),
"n_questions": len(record.get("questions") or []),
}
for record in records
]
result = {
**report,
"eval_run_id": eval_run_id,
**summary,
"per_sample": per_sample,
"artifacts": str(dest),
}
if summary.get("score_challenger") is not None and summary.get("score_king") is not None:
result["challenger_won"] = challenger_beats_king(
float(summary["score_challenger"]), float(summary["score_king"])
)
print(json.dumps({k: result[k] for k in result if k not in {"preview"}}, indent=2), flush=True)
print(f"artifacts: {dest}/official-judge.json", flush=True)
return result
def latest_run(runs_dir: Path = DEFAULT_RUNS_DIR) -> Path:
runs = sorted(
(p for p in Path(runs_dir).iterdir() if (p / "generated-samples.jsonl").is_file()),
key=lambda p: p.name,
reverse=True,
)
if not runs:
raise SystemExit(f"no duel runs with generated-samples.jsonl under {runs_dir}")
return runs[0]