dev-strender's picture
Replace v24-era demo with v34 pipeline demo (engine-vendored bundle)
9c84f9d verified
Raw History Blame Contribute Delete
5.42 kB
"""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}