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}