ValueArena / app.py
invi-bhagyesh's picture
Upload 14 files
86974b4 verified
Raw History Blame Contribute Delete
20.4 kB
"""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()