"""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}