"""EigenBench Pipeline Space — train BTD + bootstrap, upload to ValueArena.""" from __future__ import annotations import matplotlib matplotlib.use("Agg") import json import math import os import re import shutil import tempfile from datetime import datetime, timezone from pathlib import Path import gradio as gr import numpy as np import torch from sklearn.model_selection import train_test_split from torch import nn from torch.utils.data import DataLoader from pipeline.utils import ( load_records, extract_comparisons_with_ties_criteria, handle_inconsistencies_with_ties_criteria, ) from pipeline.train import ( CriteriaComparisons, Comparisons, CriteriaVectorBTD, VectorBT, build_model_labels, eigentrust_to_elo, train_vector_bt, group_split_comparisons, save_uv_embedding_plot, save_eigenbench_plot, run_bootstrap, ) from pipeline.trust import ( compute_trust_matrix, compute_trust_matrix_ties, eigentrust, row_normalize, ) HF_REPO = "invi-bhagyesh/ValueArena" GIT_REPO = "https://github.com/jchang153/EigenBench" # ── Spec parsing (from upload_results.py) ── def parse_spec(spec_content: str) -> dict: namespace = {"min": min, "max": max, "bool": bool, "True": True, "False": False} exec(spec_content, namespace) return namespace["RUN_SPEC"] def detect_model_type(model_id: str) -> dict: if model_id.startswith("hf_local:"): hf_path = model_id[len("hf_local:"):] parts = hf_path.split("/") if len(parts) >= 3: return {"id": model_id, "type": "lora", "base_model": "/".join(parts[:2]), "adapter": hf_path} else: return {"id": model_id, "type": "base", "base_model": hf_path, "adapter": None} else: return {"id": model_id, "type": "api", "base_model": None, "adapter": None} def build_meta(name, spec, log, et_scores, git_commit, git_repo): models = {} for model_name, model_id in spec.get("models", {}).items(): models[model_name] = detect_model_type(model_id) return { "name": name, "timestamp": datetime.now(timezone.utc).isoformat(), "git_commit": git_commit, "git_repo": git_repo, "models": models, "dataset": spec.get("dataset", {}), "constitution": spec.get("constitution", {}), "training": {k: v for k, v in spec.get("training", {}).items() if k not in ("bootstrap", "enabled")}, "collection": {k: v for k, v in spec.get("collection", {}).items() if k not in ("enabled", "evaluations_path", "cached_responses_path")}, "bootstrap": spec.get("training", {}).get("bootstrap", {}), "log": log, "eigentrust": et_scores, } def build_summary_from_eigentrust(et_scores, model_names): n = len(model_names) rows = [] for i, name in enumerate(model_names): trust = et_scores[i] if i < len(et_scores) else 0.0 elo = 1500.0 + 400.0 * math.log10(max(n * trust, 1e-12)) rows.append({"model_index": i, "model_name": name, "elo_mean": elo, "elo_std": 0.0, "elo_ci_lower": elo, "elo_ci_upper": elo}) rows.sort(key=lambda r: r["elo_mean"], reverse=True) return rows def parse_log_train(log_path): result = {} with open(log_path) as f: for line in f: line = line.strip() if "=" in line: key, val = line.split("=", 1) key, val = key.strip(), val.strip() try: result[key] = int(val) except ValueError: try: result[key] = float(val) except ValueError: result[key] = val return result def parse_eigentrust_file(et_path): text = Path(et_path).read_text() numbers = re.findall(r"[\d.]+(?:e[+-]?\d+)?", text) return [float(x) for x in numbers] def build_index_entry(name, meta, summary_path, group=None, note=None): with open(summary_path) as f: summary = json.load(f) top = summary[0] if summary else {} constitution_path = meta.get("constitution", {}).get("path", "") constitution_name = Path(constitution_path).stem if constitution_path else "" dataset_path = meta.get("dataset", {}).get("path", "") scenario_name = Path(dataset_path).stem if dataset_path else "" ds = meta.get("dataset", {}) start = ds.get("start", 0) count = ds.get("count", 0) scenario_range = f"{scenario_name} [{start}-{start + count}]" if scenario_name else "" return { "slug": name, "name": name, "group": group, "note": note, "timestamp": meta["timestamp"], "git_commit": meta.get("git_commit"), "models_count": len(meta.get("models", {})), "constitution": constitution_name, "scenario": scenario_range, "sampler_mode": meta.get("collection", {}).get("sampler_mode"), "btd_model": meta.get("training", {}).get("model"), "dims": meta.get("training", {}).get("dims"), "top_model": top.get("model_name", ""), "top_elo": round(top.get("elo_mean", 0), 1), "test_loss": meta.get("log", {}).get("test_loss"), } # ── Pipeline orchestration ── def run_full_pipeline(eval_path, spec, out_root, progress): train_cfg = spec.get("training", {}) constitution_cfg = spec.get("constitution", {}) num_criteria = int(constitution_cfg.get("num_criteria", 1)) progress(0.05, desc="Loading evaluations...") data = load_records(str(eval_path)) if not data: raise gr.Error("evaluations.jsonl is empty or could not be parsed") progress(0.10, desc="Extracting comparisons...") comparisons, _, extracted_name_map = extract_comparisons_with_ties_criteria( data, num_criteria=num_criteria, verbose=True, return_name_map=True, ) comparisons = handle_inconsistencies_with_ties_criteria(comparisons) if not train_cfg.get("separate_criteria", False): comparisons = [[0] + i[1:] for i in comparisons] if not comparisons: raise gr.Error("No valid comparisons extracted from evaluations") num_models = len(set([i[2] for i in comparisons] + [i[3] for i in comparisons] + [i[4] for i in comparisons])) num_criteria_eff = len(set([i[0] for i in comparisons])) model_labels = build_model_labels(num_models, spec.get("models", {}), extracted_name_map) model_kind = train_cfg.get("model", "btd_ties") batch_size = int(train_cfg.get("batch_size", 32)) lr = float(train_cfg.get("lr", 1e-3)) weight_decay = float(train_cfg.get("weight_decay", 0.0)) max_epochs = int(train_cfg.get("max_epochs", 1000)) device = "cpu" dims = list(train_cfg.get("dims", [2])) # Train/test split if train_cfg.get("group_split", False): train_comps, test_comps = group_split_comparisons( comparisons, test_size=float(train_cfg.get("test_size", 0.2)), random_state=42, verbose=True, ) else: train_comps, test_comps = train_test_split( comparisons, test_size=float(train_cfg.get("test_size", 0.2)), random_state=42, shuffle=True, ) result = {} for d in dims: out_dir = os.path.join(out_root, f"btd_d{d}") os.makedirs(out_dir, exist_ok=True) progress(0.15, desc=f"Training BTD (dim={d})...") if model_kind == "btd_ties": model = CriteriaVectorBTD(num_criteria_eff, num_models, d) train_loader = DataLoader(CriteriaComparisons(train_comps), batch_size=batch_size, shuffle=True) test_loader = DataLoader(CriteriaComparisons(test_comps), batch_size=batch_size, shuffle=False) criterion_mode, use_btd = True, True else: model = VectorBT(num_models, d) train_loader = DataLoader(Comparisons([[0] + c[1:] for c in train_comps]), batch_size=batch_size, shuffle=True) test_loader = DataLoader(Comparisons([[0] + c[1:] for c in test_comps]), batch_size=batch_size, shuffle=False) criterion_mode, use_btd = False, False loss_history = train_vector_bt( model=model, dataloader=train_loader, lr=lr, weight_decay=weight_decay, max_epochs=max_epochs, device=device, save_path=out_dir, normalize=False, use_btd=use_btd, criterion_mode=criterion_mode, verbose=True, ) progress(0.30, desc="Computing test loss...") model.eval() loss_fn = nn.CrossEntropyLoss() if use_btd else nn.BCELoss() total_test_loss = 0.0 with torch.no_grad(): for batch in test_loader: if criterion_mode: c, i, j, k, r = batch c, i, j, k = c.to(device), i.to(device), j.to(device), k.to(device) else: i, j, k, r = batch i, j, k = i.to(device), j.to(device), k.to(device) r = r.to(device) if use_btd: r = r.long() logits = model(c, i, j, k) loss = loss_fn(logits, r) else: p = model(i, j, k) loss = loss_fn(p, r) total_test_loss += loss.item() * r.size(0) avg_test_loss = total_test_loss / max(1, len(test_loader.dataset)) progress(0.35, desc="Computing EigenTrust...") if use_btd: T = compute_trust_matrix_ties(model, device) t = eigentrust(T, alpha=0, verbose=True) else: S = compute_trust_matrix(model, device) C = row_normalize(S) t = eigentrust(C, alpha=0, verbose=True) elo_np = eigentrust_to_elo(t.detach().cpu().numpy(), num_models) et_list = t.detach().cpu().numpy().tolist() progress(0.38, desc="Generating plots...") try: save_uv_embedding_plot(model=model, model_names=model_labels, save_path=os.path.join(out_dir, "uv_embeddings_pca.png")) except Exception as e: print(f"Skipping u/v PCA plot: {e}") try: save_eigenbench_plot(model_names=model_labels, eigentrust_elo=elo_np, save_path=os.path.join(out_dir, "eigenbench.png")) except Exception as e: print(f"Skipping eigenbench plot: {e}") # Write eigentrust.txt with open(os.path.join(out_dir, "eigentrust.txt"), "w") as f: f.write("EigenTrust scores:\n") f.write(np.array2string(t.cpu().numpy(), separator=", ")) f.write("\n") # Write log_train.txt with open(os.path.join(out_dir, "log_train.txt"), "w") as f: f.write(f"train_datasize = {len(train_comps)}\n") f.write(f"test_datasize = {len(test_comps)}\n") f.write(f"num_models = {num_models}\n") f.write(f"num_criteria = {num_criteria_eff}\n") f.write(f"dim = {d}\n") f.write(f"lr = {lr}\n") f.write(f"epochs = {max_epochs}\n") f.write(f"min_train_loss = {np.round(min(loss_history), 6)}\n") f.write(f"test_loss = {np.round(avg_test_loss, 6)}\n") # Bootstrap bootstrap_cfg = train_cfg.get("bootstrap", {}) has_bootstrap = bootstrap_cfg and bootstrap_cfg.get("enabled", False) if has_bootstrap: n_bootstraps = int(bootstrap_cfg.get("n_bootstraps", 100)) bootstrap_dir = os.path.join(out_dir, "bootstrap") def _progress_fn(current, total): frac = 0.40 + 0.50 * (current / total) progress(frac, desc=f"Bootstrap {current}/{total}...") progress(0.40, desc=f"Bootstrap resampling (0/{n_bootstraps})...") run_bootstrap( comparisons=comparisons, num_models=num_models, num_criteria=num_criteria_eff, model_kind=model_kind, dim=d, model_labels=model_labels, output_dir=bootstrap_dir, n_bootstraps=n_bootstraps, random_seed=int(bootstrap_cfg.get("random_seed", 42)), batch_size=batch_size, lr=lr, weight_decay=weight_decay, max_epochs=max_epochs, device=device, save_models=False, save_trust_matrices=bool(bootstrap_cfg.get("save_trust_matrices", True)), verbose=True, progress_fn=_progress_fn, ) # Store result for first dim (used for upload) if not result: result = { "btd_dir": Path(out_dir), "et_list": et_list, "model_labels": model_labels, "has_bootstrap": has_bootstrap, "log": parse_log_train(os.path.join(out_dir, "log_train.txt")), } return result # ── Upload to ValueArena ── def stage_and_upload(run_name, result, tmpdir, spec, git_commit, group, note, eval_path, progress): from huggingface_hub import HfApi, CommitOperationAdd, hf_hub_download progress(0.92, desc="Staging files for upload...") btd_dir = result["btd_dir"] et_list = result["et_list"] log = result["log"] git_repo = GIT_REPO meta = build_meta(run_name, spec, log, et_list, git_commit, git_repo) # Stage staging = Path(tmpdir) / "staging" dest = staging / "runs" / run_name dest.mkdir(parents=True, exist_ok=True) (dest / "images").mkdir(exist_ok=True) with open(dest / "meta.json", "w") as f: json.dump(meta, f, indent=2) # Summary summary_path = btd_dir / "bootstrap" / "summary.json" if summary_path.exists(): shutil.copy2(summary_path, dest / "summary.json") else: summary_data = build_summary_from_eigentrust(et_list, result["model_labels"]) with open(dest / "summary.json", "w") as f: json.dump(summary_data, f, indent=2) summary_path = dest / "summary.json" # Images for img_name, src in [ ("eigenbench.png", btd_dir / "eigenbench.png"), ("training_loss.png", btd_dir / "training_loss.png"), ("uv_embeddings_pca.png", btd_dir / "uv_embeddings_pca.png"), ("bootstrap_elo.png", btd_dir / "bootstrap" / "bootstrap_elo.png"), ]: if src.exists(): shutil.copy2(src, dest / "images" / img_name) # Evaluations if eval_path and Path(eval_path).exists(): shutil.copy2(eval_path, dest / "evaluations.jsonl") progress(0.95, desc="Uploading to HuggingFace...") token = os.environ.get("HF_TOKEN") api = HfApi(token=token) # Upload staged files staged_files = sorted(f for f in staging.rglob("*") if f.is_file()) operations = [] for fpath in staged_files: rel = fpath.relative_to(staging) operations.append(CommitOperationAdd(path_in_repo=str(rel), path_or_fileobj=fpath.read_bytes())) # Update index try: index_path = hf_hub_download(repo_id=HF_REPO, filename="index.json", repo_type="dataset", token=token) with open(index_path) as f: index = json.load(f) except Exception: index = {"last_updated": None, "runs": []} entry = build_index_entry(run_name, meta, dest / "summary.json", group=group, note=note) index["runs"] = [r for r in index["runs"] if r["slug"] != run_name] index["runs"].append(entry) index["runs"].sort(key=lambda r: r.get("timestamp", ""), reverse=True) index["last_updated"] = datetime.now(timezone.utc).isoformat() index_bytes = json.dumps(index, indent=2).encode() operations.append(CommitOperationAdd(path_in_repo="index.json", path_or_fileobj=index_bytes)) api.create_commit( repo_id=HF_REPO, repo_type="dataset", operations=operations, commit_message=f"Add run: {run_name}", ) progress(1.0, desc="Done!") return meta # ── Gradio handler ── def on_run(secret_str, evaluations_file, spec_file, run_name_str, group_str, note_str, git_commit_str, progress=gr.Progress()): expected = os.environ.get("SPACE_SECRET", "") if not expected or secret_str != expected: raise gr.Error("Invalid secret") if evaluations_file is None: raise gr.Error("Please upload evaluations.jsonl") if spec_file is None: raise gr.Error("Please upload spec.py") if not run_name_str or not run_name_str.strip(): raise gr.Error("Run name is required") run_name = run_name_str.strip() group = group_str.strip() or None if group_str else None note = note_str.strip() or None if note_str else None git_commit = git_commit_str.strip() or None if git_commit_str else None # Read spec spec_content = Path(spec_file).read_text() spec = parse_spec(spec_content) with tempfile.TemporaryDirectory() as tmpdir: # Copy evaluations to temp dir eval_path = Path(tmpdir) / "evaluations.jsonl" shutil.copy2(evaluations_file, eval_path) # Run pipeline result = run_full_pipeline(eval_path, spec, tmpdir, progress) # Upload meta = stage_and_upload(run_name, result, tmpdir, spec, git_commit, group, note, eval_path, progress) # Copy images out of tmpdir so Gradio can serve them import uuid persist_dir = Path("/tmp") / f"va_{uuid.uuid4().hex[:8]}" persist_dir.mkdir(exist_ok=True) btd_dir = result["btd_dir"] imgs = {} for key, path in [ ("eigenbench", btd_dir / "eigenbench.png"), ("bootstrap", btd_dir / "bootstrap" / "bootstrap_elo.png"), ("uv", btd_dir / "uv_embeddings_pca.png"), ("loss", btd_dir / "training_loss.png"), ]: if path.exists(): dest = persist_dir / f"{key}.png" shutil.copy2(path, dest) imgs[key] = str(dest) else: imgs[key] = None # Read summary summary_path = btd_dir / "bootstrap" / "summary.json" if not summary_path.exists(): summary_path = Path(tmpdir) / "staging" / "runs" / run_name / "summary.json" summary = json.loads(summary_path.read_text()) if summary_path.exists() else [] return ( f"Uploaded '{run_name}' to ValueArena\nhttps://valuearena.github.io/run.html?run={run_name}", imgs["eigenbench"], imgs["bootstrap"], imgs["uv"], imgs["loss"], meta, summary, ) # ── Gradio UI ── with gr.Blocks(title="EigenBench", theme=gr.themes.Soft()) as demo: gr.Markdown("# EigenBench Pipeline\nUpload evaluations + spec, run BTD training + bootstrap, publish to [ValueArena](https://valuearena.github.io).") secret = gr.Textbox(label="Secret", type="password", placeholder="Required") with gr.Row(): with gr.Column(scale=2): evaluations_file = gr.File(label="evaluations.jsonl", file_types=[".jsonl", ".json"]) spec_file = gr.File(label="spec.py", file_types=[".py"]) with gr.Column(scale=1): run_name = gr.Textbox(label="Run Name", placeholder="e.g., oct/goodness") group = gr.Textbox(label="Group (optional)", placeholder="e.g., oct") note = gr.Textbox(label="Note (optional)", placeholder="e.g., 12 persona LoRAs") git_commit = gr.Textbox(label="Git Commit (optional)", placeholder="abc1234...") run_btn = gr.Button("Run Pipeline & Upload", variant="primary") status = gr.Textbox(label="Status", interactive=False, lines=2) gr.Markdown("### Results") with gr.Row(): eigenbench_img = gr.Image(label="EigenBench Elo", type="filepath") bootstrap_img = gr.Image(label="Bootstrap CI", type="filepath") with gr.Row(): uv_img = gr.Image(label="UV Embeddings PCA", type="filepath") loss_img = gr.Image(label="Training Loss", type="filepath") with gr.Accordion("meta.json", open=False): meta_json = gr.JSON(label="meta.json") with gr.Accordion("Summary (Elo rankings)", open=False): summary_json = gr.JSON(label="summary.json") run_btn.click( fn=on_run, inputs=[secret, evaluations_file, spec_file, run_name, group, note, git_commit], outputs=[status, eigenbench_img, bootstrap_img, uv_img, loss_img, meta_json, summary_json], ) demo.queue(max_size=30, default_concurrency_limit=1) demo.launch()