Spaces:
Running
Running
File size: 5,424 Bytes
9c84f9d | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 | """Sync run artifacts to MongoDB for dashboard consumption."""
import json
import logging
from pathlib import Path
import click
logger = logging.getLogger(__name__)
@click.command("sync")
@click.argument("project")
@click.option(
"--run", "-r", "run_id", default=None, help="Specific run ID to sync (default: sync all)"
)
@click.option(
"--mongo-uri",
envvar="MONGODB_URI",
default="mongodb://localhost:27017",
help="MongoDB connection URI",
)
@click.option("--db", "database", default="solar_eval", help="MongoDB database name")
@click.pass_context
def sync_to_mongodb(
ctx: click.Context, project: str, run_id: str | None, mongo_uri: str, database: str
) -> None:
"""Sync run artifacts to MongoDB for the dashboard."""
from solar_eval.stores.mongodb import MongoDBStore
cfg = ctx.obj["config"]
artifacts_dir = cfg.artifacts_dir(project)
if not artifacts_dir.exists():
raise click.ClickException(f"Artifacts directory not found: {artifacts_dir}")
store = MongoDBStore(connection_uri=mongo_uri, database=database)
store.ensure_indexes()
if run_id:
run_dirs = [artifacts_dir / run_id]
if not run_dirs[0].exists():
raise click.ClickException(f"Run not found: {run_dirs[0]}")
else:
run_dirs = sorted(
[d for d in artifacts_dir.iterdir() if d.is_dir() and (d / "run.json").exists()]
)
if not run_dirs:
click.secho("No runs found to sync.", fg="yellow")
return
click.secho(f"Syncing {len(run_dirs)} run(s) to MongoDB ({mongo_uri}/{database})...", fg="blue")
synced = 0
skipped = 0
for run_dir in run_dirs:
try:
counts = _sync_single_run(store, project, run_dir, cfg)
synced += 1
click.secho(
f" {run_dir.name}: {counts['samples']} samples, {counts['diffs']} diffs",
fg="green",
)
except Exception as e:
skipped += 1
click.secho(f" {run_dir.name}: SKIP ({e})", fg="yellow")
store.close()
click.secho(f"\nDone: {synced} synced, {skipped} skipped.", fg="green", bold=True)
def _sync_single_run(store, project: str, run_dir: Path, cfg) -> dict[str, int]:
"""Sync a single run directory to MongoDB. Returns counts."""
run_json = run_dir / "run.json"
run_meta = json.loads(run_json.read_text())
task = run_meta.get("task", "")
run_id = run_dir.name
# Build experiment document
experiment = {
"project": project,
"run_id": run_id,
"task": task,
"model": run_meta.get("model", ""),
"temperature": run_meta.get("temperature", 0.0),
"prompt_version": run_meta.get("prompt_version", ""),
"status": run_meta.get("status", "unknown"),
"total_samples": run_meta.get("total_samples", 0),
"completed_samples": run_meta.get("completed_samples", 0),
"started_at": run_meta.get("started_at", ""),
"completed_at": run_meta.get("completed_at", ""),
}
# Enrich with evaluation scores
eval_file = run_dir / "evaluation.json"
if eval_file.exists():
eval_data = json.loads(eval_file.read_text())
experiment["scores"] = eval_data.get("scores", {})
experiment["overall_score"] = eval_data.get("overall_score", 0.0)
# Enrich with prompt/pipeline metadata from run.json
experiment["prompt_id"] = run_meta.get("prompt", run_meta.get("prompt_version", ""))
experiment["pipeline_id"] = run_meta.get("pipeline", "")
experiment["reasoning_effort"] = run_meta.get("reasoning_effort", "")
experiment_id = store.upsert_experiment(experiment)
# Sync sample results
sample_count = 0
results_file = run_dir / "results.jsonl"
if results_file.exists():
samples = []
for line in results_file.read_text().strip().split("\n"):
if not line:
continue
result = json.loads(line)
samples.append(
{
"experiment_id": experiment_id,
"sample_idx": result.get("sample_idx", 0),
"input": result.get("input", {}),
"output": result.get("output", ""),
"golden": result.get("golden", ""),
"input_tokens": result.get("input_tokens", 0),
"output_tokens": result.get("output_tokens", 0),
"inference_time_ms": result.get("inference_time_ms", 0),
}
)
sample_count = store.bulk_insert_samples(samples)
# Sync eval details
diff_count = 0
eval_details_file = run_dir / "eval_details.jsonl"
if eval_details_file.exists():
diffs = []
for line in eval_details_file.read_text().strip().split("\n"):
if not line:
continue
detail = json.loads(line)
diffs.append(
{
"experiment_id": experiment_id,
"sample_idx": detail.get("sample_idx", 0),
"score": detail.get("score", 0.0),
"category_scores": detail.get("category_scores", {}),
"error_counts": detail.get("error_counts", {}),
}
)
diff_count = store.bulk_insert_diffs(diffs)
return {"samples": sample_count, "diffs": diff_count}
|