"""Run management commands — local execution.""" import asyncio import json import subprocess import time from datetime import datetime from pathlib import Path from typing import Any import click from rich.progress import Progress, SpinnerColumn, BarColumn, TextColumn, TimeElapsedColumn from solar_eval.cli.commands.projects import _resolve_project_id from solar_eval.cli.formatters import format_runs_table, format_results_table, format_scores #: run 이 끝난 뒤 자동 적재에 허용하는 시간. 전체 아카이브가 아니라 프로젝트 하나만 #: 적재하므로 보통 수 초면 끝난다 — 넘어가면 뭔가 잘못된 것이라 기다리지 않는다. INGEST_TIMEOUT_SECONDS = 180 @click.group("runs") def runs_group() -> None: """Manage evaluation runs.""" @runs_group.command("list") @click.argument("project") @click.option("--task", "-t", default=None, help="Filter by task name") @click.pass_context def list_runs(ctx: click.Context, project: str, task: str | None) -> None: """List runs for a project.""" cfg = ctx.obj["config"] if not cfg.is_remote: # List local run directories artifacts = cfg.artifacts_dir(project) if not artifacts.exists(): click.secho("No local runs found.", fg="yellow") return runs = [] for d in sorted(artifacts.iterdir()): if d.is_dir() and (d / "run.json").exists(): run_data = json.loads((d / "run.json").read_text()) if task and run_data.get("task") != task: continue runs.append( { "id": d.name, "task": run_data.get("task", "-"), "model": run_data.get("model", "-"), "prompt": run_data.get("prompt", run_data.get("prompt_version", "-")), "pipeline": run_data.get("pipeline", "-"), "status": run_data.get("status", "unknown"), "total_samples": run_data.get("total_samples", 0), "completed_samples": run_data.get("completed_samples", 0), "created_at": run_data.get("started_at", "-"), } ) if not runs: click.secho("No local runs found.", fg="yellow") return click.secho(format_runs_table(runs), fg="blue") return from solar_eval.cli.client import EvalClient client = EvalClient(cfg.remote_url, cfg.timeout) project_id = _resolve_project_id(client, project) params = {"task": task} if task else {} runs = client.get(f"/api/projects/{project_id}/runs", params=params) if not runs: click.secho("No runs found.", fg="yellow") return click.secho(format_runs_table(runs), fg="blue") @runs_group.command("start") @click.argument("project") @click.option("--task", "-t", required=True, help="Task name") @click.option("--prompt", "-p", "prompt_name", required=True, help="Prompt name (YAML in prompts/)") @click.option("--model", "-m", required=True, help="Model name (e.g. solar-pro3)") @click.option( "--pipeline", "-P", "pipeline_name", default=None, help="Pipeline template name (overrides project default)", ) @click.option( "--reasoning-effort", "-re", default=None, help="Reasoning effort (low/medium/high)", ) @click.option("--temperature", "-T", default=0.0, type=float, help="Temperature (default: 0.0)") @click.option("--max-tokens", default=4000, type=int, help="Max tokens (default: 4000)") @click.option("--max-workers", "-w", default=5, type=int, help="Max parallel workers") @click.option( "--limit", "-n", default=None, type=int, help="Limit number of samples (for quick testing)" ) @click.option( "--note", default=None, help="이 run 을 왜 돌리는지 한 줄 (run.json 에 기록되고 evalhub 가 그대로 읽는다)", ) @click.option( "--set", "set_options", multiple=True, metavar="PATH=VALUE", help=( "파이프라인 설정 한 축만 바꿔서 돌린다 (파일을 만들지 않는다). " "예: --set steps.basic_correction.reasoning_effort=high. 여러 번 쓸 수 있다." ), ) @click.option( "--no-ingest", is_flag=True, help="run 이 끝난 뒤 evalhub 자동 적재를 건너뛴다 (여러 run 을 돌리고 한 번에 적재할 때)", ) @click.pass_context def start_run( ctx: click.Context, project: str, task: str, prompt_name: str, model: str, pipeline_name: str | None, reasoning_effort: str | None, temperature: float, max_tokens: int, max_workers: int, limit: int | None, note: str | None, set_options: tuple[str, ...], no_ingest: bool, ) -> None: """Start an evaluation run. Example: solar-eval runs start chosun-proofreading -t ci-v1-all \\ -p dev_260329_pro3_v1 -m solar-pro3 -P 251231_pro3 \\ --note "v24 SC+judge 가 ci-v1 에서도 prod 를 넘는지" solar-eval runs start chosun-proofreading -t paragraph-small \\ -p prompt_dev_v23 -m solar-pro3 -P pipeline_dev_v24 \\ --set steps.basic_correction.repeat_temperatures='[0.0, 0.5]' \\ --note "SC 온도 폭 스윕" """ cfg = ctx.obj["config"] _start_local( cfg, project=project, task=task, prompt_name=prompt_name, model=model, pipeline_name=pipeline_name, reasoning_effort=reasoning_effort, temperature=temperature, max_tokens=max_tokens, max_workers=max_workers, limit=limit, note=note, set_options=set_options, no_ingest=no_ingest, ) def _resolve_prompt(prompts_dir: Path, prompt_name: str) -> Path: """Resolve prompt from prompts/ directory (directory or YAML file). Resolution order: 1. prompts/{name}/ (directory TXT format) 2. prompts/{name}.yaml (flat YAML layout) 3. prompts/_library/{name}.yaml (legacy layout) """ # Directory TXT format (new) dir_path = prompts_dir / prompt_name if dir_path.is_dir(): return dir_path # Flat YAML layout yaml_path = prompts_dir / f"{prompt_name}.yaml" if yaml_path.exists(): return yaml_path # Legacy _library layout legacy_path = prompts_dir / "_library" / f"{prompt_name}.yaml" if legacy_path.exists(): return legacy_path raise click.ClickException( f"Prompt '{prompt_name}' not found. Searched:\n" f" {dir_path}/ (directory)\n" f" {yaml_path} (YAML)\n" f" {legacy_path} (legacy)" ) def _resolve_pipeline_file(project_dir: Path, pipeline_name: str) -> Path: """Resolve pipeline YAML from pipelines/ directory.""" path = project_dir / "pipelines" / f"{pipeline_name}.yaml" if not path.exists(): raise click.ClickException(f"Pipeline '{pipeline_name}' not found: {path}") return path def _compose_pipeline_config(project_dir: Path, pipeline_name: str, overrides: dict) -> dict: """파이프라인을 `extends`/override 까지 해석해 돌려준다. 조립 실패(없는 base, 없는 스텝을 가리키는 override 등)는 run 을 시작하기 전에 죽인다 — 효과 없는 override 로 돌아간 run 은 결론을 조용히 오염시킨다. """ from solar_eval.core.pipeline_compose import PipelineCompositionError, compose_pipeline from solar_eval.core.project_loader import load_pipeline_file pipelines_dir = _resolve_pipeline_file(project_dir, pipeline_name).parent try: composed = load_pipeline_file(pipelines_dir, pipeline_name) if overrides: composed = compose_pipeline(composed, overrides=overrides) return composed except (PipelineCompositionError, FileNotFoundError) as e: raise click.ClickException(str(e)) from None def _build_run_name(task: str, pipeline_name: str, prompt_name: str) -> str: """Build a readable run name: {task}__{pipeline}__{prompt}__{YYMMDD}_{HHMMSS}""" ts = datetime.now().strftime("%y%m%d_%H%M%S") return f"{task}__{pipeline_name}__{prompt_name}__{ts}" _RUN_DIR_MAX_ATTEMPTS = 20 _RUN_DIR_RETRY_DELAY_SECONDS = 0.05 def _make_run_dir( artifacts_dir: Path, task: str, pipeline_name: str, prompt_name: str ) -> tuple[str, Path]: """run_name 을 짓고 그 디렉터리를 원자적으로 만든다. 같은 초에 같은 (task, pipeline, prompt) 조합으로 두 run 이 시작되면 `_build_run_name` 이 같은 이름을 낸다. `mkdir(exist_ok=False)` 를 쓰면 두 프로세스가 동시에 같은 이름으로 mkdir 해도 정확히 하나만 성공하고, 진 쪽은 새 타임스탬프로 재시도한다 — 파일 포맷 자체는 안 바꾸므로 evalhub 의 run_id 파싱과 호환된다. """ for _ in range(_RUN_DIR_MAX_ATTEMPTS): run_name = _build_run_name(task, pipeline_name, prompt_name) run_dir = artifacts_dir / run_name try: run_dir.mkdir(parents=True, exist_ok=False) return run_name, run_dir except FileExistsError: time.sleep(_RUN_DIR_RETRY_DELAY_SECONDS) raise click.ClickException(f"Could not allocate a unique run directory under {artifacts_dir}") def _config_source(cfg, project: str) -> str: """run 의 config 출처 — 레포 관리 프로젝트면 레포 상대 경로, 아니면 데이터 홈 절대 경로.""" config_dir = cfg.config_dir(project) if cfg.is_repo_managed(project) and cfg.repo_root is not None: return str(config_dir.relative_to(cfg.repo_root)) return str(config_dir) def _git_commit(repo_root: Path | None) -> str | None: """재현용 도장: 레포 HEAD 커밋. 작업 트리가 더러우면 '-dirty' 를 붙인다.""" if repo_root is None: return None try: head = subprocess.run( ["git", "-C", str(repo_root), "rev-parse", "--short", "HEAD"], capture_output=True, text=True, check=True, ).stdout.strip() status = subprocess.run( ["git", "-C", str(repo_root), "status", "--porcelain"], capture_output=True, text=True, check=True, ).stdout.strip() return f"{head}-dirty" if status else head except (OSError, subprocess.CalledProcessError): return None def _start_local( cfg, project: str, task: str, prompt_name: str, model: str, pipeline_name: str | None, reasoning_effort: str | None, temperature: float, max_tokens: int, max_workers: int, limit: int | None, note: str | None = None, set_options: tuple[str, ...] = (), no_ingest: bool = False, ) -> None: """Start run locally with the new prompt/model/pipeline interface.""" from solar_eval.core.dataset_loader import DatasetLoader from solar_eval.core.fingerprint import pipeline_fingerprint, prompt_fingerprint from solar_eval.core.pipeline_compose import PipelineCompositionError, parse_set_option from solar_eval.core.project_loader import load_all_project_configs from solar_eval.core.runner import BatchRunner from solar_eval.models.prompt_version import ( detect_prompt_format, load_prompt_messages, load_step_prompts, ) from solar_eval.providers import UpstageProvider from solar_eval.providers.openai_provider import OpenAIProvider from solar_eval.stores import JsonlStore # Load project config configs = load_all_project_configs(cfg.projects_dir, cfg.config_dirs) project_config = next((c for c in configs if c["name"] == project), None) if not project_config: raise click.ClickException(f"Project not found in {cfg.projects_dir}: {project}") task_config = next((t for t in project_config.get("tasks", []) if t["name"] == task), None) if not task_config: raise click.ClickException(f"Task '{task}' not found in project '{project}'") # Resolve prompt file prompts_dir = cfg.prompts_dir(project) prompt_path = _resolve_prompt(prompts_dir, prompt_name) fmt = detect_prompt_format(prompt_path) if fmt == "multi_step": step_prompts = load_step_prompts(prompt_path) prompt = { "step_prompts": step_prompts, "model": model, "temperature": temperature, "max_tokens": max_tokens, } else: messages = load_prompt_messages(prompt_path) prompt = { "messages": messages, "model": model, "temperature": temperature, "max_tokens": max_tokens, } if reasoning_effort: prompt["reasoning_effort"] = reasoning_effort # Resolve pipeline — config 는 레포 관리 프로젝트면 git 쪽이 정본. # `extends` 상속과 `--set` override 는 여기서 전부 해석돼, 실행 엔진에는 완성된 # steps 목록만 넘어간다. try: cli_overrides = parse_set_option(set_options) except PipelineCompositionError as e: raise click.ClickException(str(e)) from None effective_pipeline = pipeline_name if pipeline_name: pipeline_config = _compose_pipeline_config( cfg.config_dir(project), pipeline_name, cli_overrides ) task_config = {**task_config, "pipeline_config": pipeline_config} else: # Use default from project.yaml (already resolved by project_loader) pc = task_config.get("pipeline_config", {}) effective_pipeline = pc.get("name", "default") if isinstance(pc, dict) else str(pc) if cli_overrides: raise click.ClickException("--set requires an explicit --pipeline/-P") pipeline_config = pc if isinstance(pc, dict) else {} # Build run name and directory -- 이름 충돌은 여기서 원자적으로 걸러진다(F10). run_name, run_dir = _make_run_dir( cfg.artifacts_dir(project), task, effective_pipeline, prompt_name ) store = JsonlStore(run_dir) # Save run metadata run_meta = { "task": task, "model": model, "prompt": prompt_name, "pipeline": effective_pipeline, "temperature": temperature, "max_tokens": max_tokens, "reasoning_effort": reasoning_effort, # 왜 돌렸는지 — 예전에는 별도 랩노트(experiments/*.yaml)가 담았지만 run 과 분리돼 있어 # 참조가 끊어지고 썩었다. run 과 같은 파일에 두면 끊어질 수가 없다. "note": note, # provenance: 이 run 의 config(프롬프트·파이프라인)가 어디서 왔는지 (W&B/DVC 식 도장) "config_source": _config_source(cfg, project), "git_commit": _git_commit(cfg.repo_root) if cfg.is_repo_managed(project) else None, # 내용 지문 — 이름이 아니라 내용으로 설정을 식별한다. 이름만 다른 사본, # `extends` 로 조립한 것, `--set` 으로 즉석에서 만든 것이 모두 같은 값이면 # 같은 지문을 갖는다 (core/fingerprint.py). "prompt_sha": prompt_fingerprint(prompt_path), "pipeline_sha": pipeline_fingerprint(pipeline_config), # 파일 없이 돌린 실험이면 무엇을 바꿨는지가 여기 남는다. "pipeline_overrides": cli_overrides or None, } # run_dir 은 _make_run_dir 이 이미 원자적으로 만들었다 -- 여기서 다시 mkdir 하지 않는다. (run_dir / "run.json").write_text(json.dumps(run_meta, ensure_ascii=False, indent=2)) # Log run info click.secho(f"Prompt: {prompt_name} [{fmt}]", fg="blue") click.secho(f"Model: {model}", fg="blue") click.secho(f"Pipeline: {effective_pipeline}", fg="blue") for path, value in cli_overrides.items(): click.secho(f" override: {path} = {value!r}", fg="magenta") click.secho( f"Fingerprint: prompt {run_meta['prompt_sha']} / pipeline {run_meta['pipeline_sha']}", fg="blue", ) if reasoning_effort: click.secho(f"Reasoning effort: {reasoning_effort}", fg="blue") if limit: click.secho(f"Limit: {limit} samples", fg="yellow") click.secho(f"Starting local run: {project}/{run_name}", fg="blue") click.secho(f"Output: {run_dir}", fg="blue") # Create runner provider = UpstageProvider() judge_provider = OpenAIProvider() dataset_loader = DatasetLoader(base_dir=cfg.projects_dir) runner = BatchRunner( store=store, inference_provider=provider, judge_provider=judge_provider, dataset_loader=dataset_loader, ) # Run with progress display _run_with_progress(runner, run_name, project_config, task_config, prompt, max_workers, limit) _warn_if_incomplete(run_dir) click.echo(f"\nResults saved to: {run_dir}/") if not no_ingest: _ingest_into_evalhub(cfg, project, run_dir) def _ingest_into_evalhub(cfg, project: str, run_dir: Path) -> None: """방금 만든 run 을 evalhub 로 적재한다. **실패해도 run 은 성공이다** — 결과는 이미 디스크에 있고, 적재는 나중에 다시 돌리면 되는 멱등 작업이다. 그래서 어떤 실패도 exit code 를 바꾸지 않고 "무엇이 안 됐고 어떻게 되살리는지"만 알린다. 적재를 못 했다는 사실 자체를 조용히 넘기지도 않는다 — 그러면 evalhub 화면이 낡은 채로 남는다. 레포 밖 standalone 설치에는 evalhub 자체가 없으므로 조용히 건너뛴다. """ if cfg.repo_root is None: return click.secho("\nevalhub 적재 중...", fg="blue") try: result = subprocess.run( ["pnpm", "--filter", "@poc/eval-store", "ingest", "--only", project], cwd=str(cfg.repo_root), capture_output=True, text=True, timeout=INGEST_TIMEOUT_SECONDS, ) except FileNotFoundError: _warn_ingest_skipped(project, run_dir, "pnpm 을 찾을 수 없습니다") return except subprocess.TimeoutExpired: _warn_ingest_skipped(project, run_dir, f"{INGEST_TIMEOUT_SECONDS}초 안에 끝나지 않았습니다") return except OSError as e: _warn_ingest_skipped(project, run_dir, str(e)) return if result.returncode == 0: click.secho(f" evalhub 적재 완료 — {project}", fg="green") return _warn_ingest_skipped(project, run_dir, _ingest_failure_reason(result)) def _ingest_failure_reason(result: subprocess.CompletedProcess[str]) -> str: """ingest 실패 출력에서 사람이 읽을 한 줄을 뽑는다. stderr 전체는 node 스택트레이스라 그대로 보여주면 정작 원인이 묻힌다. 가장 흔한 원인(백엔드 미기동)은 연결 거부로 나타나므로 먼저 짚는다. """ output = f"{result.stderr}\n{result.stdout}" if "ECONNREFUSED" in output or "connect ECONNREFUSED" in output: return "Postgres 에 연결할 수 없습니다 (백엔드가 꺼져 있는 것 같습니다)" for line in result.stderr.splitlines(): stripped = line.strip() if stripped and not stripped.startswith("at ") and "node_modules" not in stripped: return stripped[:160] return f"ingest 가 exit {result.returncode} 로 끝났습니다" def _warn_ingest_skipped(project: str, run_dir: Path, reason: str) -> None: """적재 실패를 알리고 되살리는 명령을 그대로 실어 준다 (exit code 는 0 유지).""" click.secho(f"\n ⚠ evalhub 적재를 건너뛰었습니다 — {reason}", fg="yellow", bold=True) click.secho(f" run 결과는 그대로 있습니다: {run_dir}", fg="yellow") click.secho(" 백엔드를 띄운 뒤 아래를 실행하면 반영됩니다:", fg="yellow") click.secho(" make evalhub-up", fg="yellow") click.secho(f" make evalhub-ingest ONLY={project}", fg="yellow") def _warn_if_incomplete(run_dir: Path) -> None: """샘플이 유실된 채 끝났으면 점수 옆에 크게 알린다. 유실은 조용하다: 실패한 샘플은 `results.jsonl` 에 행 자체가 남지 않고(error 필드도 없다) 점수는 살아남은 샘플만으로 집계된다. 실제로 429 재시도 소진으로 159건 중 17건이 빠진 run 이 "완료"로 끝나며 점수를 냈고, `run.json` 의 completed/total 을 직접 대조하기 전까지 아무도 몰랐다. 부분 결과라도 지우지는 않는다 — 재현 비용이 크고, 완료율을 알고 보면 쓸모가 있다. """ try: meta = json.loads((run_dir / "run.json").read_text()) except (OSError, json.JSONDecodeError): return total, done = meta.get("total_samples"), meta.get("completed_samples") if not isinstance(total, int) or not isinstance(done, int) or done >= total: return click.secho( f"\n ⚠ 샘플 {total - done}건이 유실됐습니다 ({done}/{total} 완료). " f"점수는 완료된 샘플만으로 계산된 값이라 다른 run 과 직접 비교할 수 없습니다.", fg="yellow", bold=True, ) click.secho( " 흔한 원인은 API 429 입니다 — `-w/--max-workers` 를 낮춰 다시 돌리세요.", fg="yellow", ) def _run_with_progress( runner, run_id, project_config, task_config, prompt, max_workers, limit, completed_results: list[dict[str, Any]] | None = None, ): """Execute a run with Rich progress bar display. `completed_results` 는 resume 이 이미 끝난 샘플을 다시 돌리지 않게 하는 다리다 -- 여기서 안 넘기면 runner 가 파라미터를 갖고 있어도 소용이 없다(F2). `execute_run` 은 치명적 실패에서 raise 하므로 이 함수도 raise 한다 -- 일부러 감싸지 않는다. 위쪽 `on_progress` 의 "error" 이벤트가 이미 사람이 읽을 메시지를 찍었고, 여기서 또 감싸면 같은 내용이 두 번 나온다. """ import logging as _logging from rich.console import Console from rich.live import Live from rich.text import Text console = Console(stderr=True) use_rich = True live_ctx: Live | None = None progress_bar = Progress( SpinnerColumn(), TextColumn("[progress.description]{task.description}"), BarColumn(), TextColumn("{task.completed}/{task.total}"), TextColumn("[progress.percentage]{task.percentage:>3.0f}%"), TimeElapsedColumn(), ) try: live_ctx = Live(progress_bar, console=console, refresh_per_second=8) live_ctx.start() except Exception: use_rich = False # Redirect pipeline/runner logs above the progress bar class _LiveLogHandler(_logging.Handler): def emit(self, record): if live_ctx: live_ctx.console.print(Text(record.getMessage(), style="dim")) _log_names = ( "solar_eval.pipelines.pipeline", "solar_eval.pipelines.steps.llm", "solar_eval.core.runner", "solar_eval.providers.upstage", ) if use_rich: _live_handler = _LiveLogHandler() for name in _log_names: lg = _logging.getLogger(name) lg.addHandler(_live_handler) lg.propagate = False _logging.getLogger("httpx").setLevel(_logging.WARNING) def _cleanup_live(): if use_rich: live_ctx.stop() for name in _log_names: _logging.getLogger(name).removeHandler(_live_handler) _logging.getLogger(name).propagate = True inference_task_id = None eval_task_id = None async def on_progress(event): nonlocal inference_task_id, eval_task_id etype = event.get("type") if etype == "dataset_loaded": total = event["total"] if use_rich: inference_task_id = progress_bar.add_task("Inference", total=total) else: click.secho(f" Dataset loaded: {total} samples", fg="blue") elif etype == "inference_progress": if use_rich and inference_task_id is not None: progress_bar.update(inference_task_id, completed=event["completed"]) elif not use_rich: c, t = event["completed"], event["total"] if c % 10 == 0 or c == t: click.echo(f" Inference: {c}/{t}") elif etype == "inference_complete": if use_rich and inference_task_id is not None: progress_bar.update(inference_task_id, completed=event["total"]) else: click.secho(f" Inference complete: {event['total']} samples", fg="green") elif etype == "evaluation_start": total = event.get("total", 0) if use_rich: eval_task_id = progress_bar.add_task("Evaluation", total=total or 1) else: click.secho(" Starting evaluation...", fg="blue") elif etype == "evaluation_progress": if use_rich and eval_task_id is not None: progress_bar.update(eval_task_id, completed=event["completed"]) elif etype == "done": if use_rich and eval_task_id is not None: progress_bar.update(eval_task_id, completed=progress_bar.tasks[eval_task_id].total) _cleanup_live() click.echo() click.secho("Run completed!", fg="green", bold=True) overall = event.get("overall_score") if overall is not None: click.echo(f" Overall Score: {overall:.4f}") scores = event.get("scores", {}) if scores: click.echo(f" Scores: {format_scores(scores)}") elif etype == "error": _cleanup_live() click.echo() click.secho(f" Run failed: {event.get('error')}", fg="red", bold=True) asyncio.run( runner.execute_run( run_id=run_id, project_config=project_config, task_config=task_config, prompt=prompt, on_progress=on_progress, max_workers=max_workers, limit=limit, completed_results=completed_results, ) ) @runs_group.command("batch") @click.argument("project") @click.option("--task", "-t", required=True, help="Task name") @click.option("--prompts", required=True, help="Comma-separated prompt names") @click.option("--model", "-m", required=True, help="Model name") @click.option( "--pipeline", "-P", "pipeline_name", default=None, help="Pipeline template name (overrides project default)", ) @click.option("--reasoning-effort", "-re", default=None, help="Reasoning effort") @click.option("--max-workers", "-w", default=5, type=int, help="Max parallel workers") @click.pass_context def batch_runs(ctx, project, task, prompts, model, pipeline_name, reasoning_effort, max_workers): """Run evaluation for multiple prompts sequentially.""" cfg = ctx.obj["config"] if cfg.is_remote: raise click.ClickException("Batch is local-mode only") prompt_list = [p.strip() for p in prompts.split(",")] click.secho(f"Batch run: {project}/{task} — prompts: {prompt_list}", fg="blue", bold=True) for pname in prompt_list: click.secho(f"\n{'=' * 50}", fg="blue") click.secho(f"Running prompt: {pname}...", fg="blue", bold=True) click.secho(f"{'=' * 50}", fg="blue") try: _start_local( cfg, project=project, task=task, prompt_name=pname, model=model, pipeline_name=pipeline_name, reasoning_effort=reasoning_effort, temperature=0.0, max_tokens=4000, max_workers=max_workers, limit=None, ) except Exception as e: click.secho(f" {pname} failed: {e}", fg="red") click.secho(f"\n{'=' * 50}", fg="blue") click.secho("Batch complete!", fg="green", bold=True) @runs_group.command("resume") @click.argument("project") @click.option("--run-id", "-r", required=True, help="Run ID to resume") @click.option("--max-workers", "-w", default=5, type=int, help="Max parallel workers") @click.pass_context def resume_run(ctx: click.Context, project: str, run_id: str, max_workers: int) -> None: """Resume an interrupted evaluation run.""" cfg = ctx.obj["config"] if cfg.is_remote: raise click.ClickException("Resume is local-mode only") from solar_eval.core.dataset_loader import DatasetLoader from solar_eval.core.project_loader import load_all_project_configs from solar_eval.core.runner import BatchRunner from solar_eval.providers import UpstageProvider from solar_eval.providers.openai_provider import OpenAIProvider from solar_eval.stores import JsonlStore # Find the run directory run_dir = cfg.artifacts_dir(project) / run_id if not run_dir.exists(): raise click.ClickException(f"Run directory not found: {run_dir}") store = JsonlStore(run_dir) run_meta = store.load_run_meta() if not run_meta: raise click.ClickException(f"No run.json found in {run_dir}") status = run_meta.get("status", "unknown") if status == "completed": raise click.ClickException(f"Run {run_id} is already completed. Nothing to resume.") task = run_meta.get("task") if not task: raise click.ClickException("Cannot determine task from run.json") # Load existing results existing_results = store.load_existing_results() completed_indices = {r["sample_idx"] for r in existing_results} click.secho(f"Resuming run: {run_id}", fg="blue", bold=True) click.secho(f" Task: {task}, Already completed: {len(completed_indices)} samples", fg="blue") # Load project/task config configs = load_all_project_configs(cfg.projects_dir, cfg.config_dirs) project_config = next((c for c in configs if c["name"] == project), None) if not project_config: raise click.ClickException(f"Project not found: {project}") task_config = next((t for t in project_config.get("tasks", []) if t["name"] == task), None) if not task_config: raise click.ClickException(f"Task '{task}' not found in project '{project}'") # Resolve prompt from run.json prompt_name = run_meta.get("prompt") pipeline_name = run_meta.get("pipeline") model = run_meta.get("model", "solar-pro2") if not prompt_name: raise click.ClickException( "Cannot determine prompt from run.json. Only runs created with the new CLI can be resumed." ) from solar_eval.models.prompt_version import ( detect_prompt_format, load_prompt_messages, load_step_prompts, ) prompt_path = _resolve_prompt(cfg.prompts_dir(project), prompt_name) fmt = detect_prompt_format(prompt_path) if fmt == "multi_step": prompt = { "step_prompts": load_step_prompts(prompt_path), "model": model, "temperature": run_meta.get("temperature", 0.0), "max_tokens": run_meta.get("max_tokens", 4000), } else: prompt = { "messages": load_prompt_messages(prompt_path), "model": model, "temperature": run_meta.get("temperature", 0.0), "max_tokens": run_meta.get("max_tokens", 4000), } if run_meta.get("reasoning_effort"): prompt["reasoning_effort"] = run_meta["reasoning_effort"] # Override pipeline if specified — 원래 run 이 `--set` 으로 바꿔 돌린 것이면 # 그 override 까지 복원해야 같은 설정으로 이어진다. if pipeline_name and pipeline_name != "default": task_config = { **task_config, "pipeline_config": _compose_pipeline_config( cfg.config_dir(project), pipeline_name, run_meta.get("pipeline_overrides") or {}, ), } click.secho(f" Model: {model}", fg="blue") # 이전 평가를 미리 지우지 않는다(F3) -- 재평가가 도중에 죽으면 복구할 점수가 # 없어진다. create_evaluation() 이 원자적 교체라 지우지 않아도 안전하다. # Setup providers and runner provider = UpstageProvider() judge_provider = OpenAIProvider() dataset_loader = DatasetLoader(base_dir=cfg.projects_dir) runner = BatchRunner( store=store, inference_provider=provider, judge_provider=judge_provider, dataset_loader=dataset_loader, ) _run_with_progress( runner, run_id, project_config, task_config, prompt, max_workers, limit=None, # 이미 끝난 샘플은 건너뛴다 -- 안 넘기면 resume 이 전 샘플을 다시 유료 호출하고 # results.jsonl 에 같은 sample_idx 를 중복 append 한다(F2). completed_results=existing_results, ) click.echo(f"\nResults saved to: {run_dir}/") @runs_group.command("eval") @click.argument("project") @click.option( "--run-id", "-r", required=True, help="Run ID to evaluate", ) @click.pass_context def eval_run(ctx: click.Context, project: str, run_id: str) -> None: """Run evaluation only on existing inference results (no new inference).""" cfg = ctx.obj["config"] if cfg.is_remote: raise click.ClickException("Eval is local-mode only") from solar_eval.core.project_loader import load_all_project_configs from solar_eval.core.runner import build_eval_sample_from_result from solar_eval.evaluators.registry import create_evaluator from solar_eval.stores import JsonlStore # Find run directory run_dir = cfg.artifacts_dir(project) / run_id if not run_dir.exists(): raise click.ClickException(f"Run directory not found: {run_dir}") store = JsonlStore(run_dir) run_meta = store.load_run_meta() if not run_meta: raise click.ClickException(f"No run.json found in {run_dir}") task = run_meta.get("task") if not task: raise click.ClickException("Cannot determine task from run.json") # Load results results = store.load_existing_results() if not results: raise click.ClickException(f"No inference results found in {run_dir}/results.jsonl") click.secho(f"Evaluating run: {run_id}", fg="blue", bold=True) click.secho(f" Task: {task}, Samples: {len(results)}", fg="blue") # Load task config for evaluator settings configs = load_all_project_configs(cfg.projects_dir, cfg.config_dirs) project_config = next((c for c in configs if c["name"] == project), None) if not project_config: raise click.ClickException(f"Project not found: {project}") task_config = next((t for t in project_config.get("tasks", []) if t["name"] == task), None) if not task_config: raise click.ClickException(f"Task '{task}' not found in project '{project}'") # 이전 평가를 미리 지우지 않는다(F3) -- 아래 재평가가 실패해도 옛 점수는 남는다. # Run evaluation evaluator = create_evaluator(task_config.get("evaluator", {"type": "llm_judge"})) # judge provider 를 안 넘기면 llm_judge 계열은 전 샘플에서 실패한다 -- 그리고 # 그 실패가 예전엔 0점 placeholder 로 덮여 조용히 completed 로 끝났다(F8). from solar_eval.providers.openai_provider import OpenAIProvider judge_provider = OpenAIProvider() from datetime import timezone async def run_eval(): # 채점 실패는 점수가 아니라 실패로 남긴다 -- aggregate() 에는 성공분만 넘긴다. eval_results: list[dict[str, Any]] = [] eval_failures: list[dict[str, Any]] = [] for result in results: # 리스트 위치가 아니라 원본 sample_idx 가 정본이다(F6). sample_idx = result["sample_idx"] try: eval_sample = build_eval_sample_from_result(result, reference=result.get("golden")) evaluator.validate_required_fields(eval_sample) eval_result = await evaluator.evaluate(sample=eval_sample, provider=judge_provider) eval_results.append({**eval_result, "sample_idx": sample_idx}) except Exception as e: click.secho(f" Sample {sample_idx} eval failed: {e}", fg="yellow") eval_failures.append({"sample_idx": sample_idx, "error": str(e)}) # Aggregate and save aggregated = evaluator.aggregate(eval_results) # 성공한 채점이 하나도 없으면 점수 자리를 비운다 -- 0.0 은 "0점"이라는 측정값이라 # "측정이 없었다"와 구분되지 않는다(runner.py 와 같은 규칙). overall_score = aggregated.get("overall_score", 0.0) if eval_results else None eval_id = await store.create_evaluation( { "run_id": run_id, "scores": aggregated.get("scores", {}), "overall_score": overall_score, "eval_model": "local", "eval_success_count": len(eval_results), "eval_failed_count": len(eval_failures), "failed_sample_indices": [f["sample_idx"] for f in eval_failures], } ) # Save per-sample eval details -- 성공/실패 둘 다 한 행씩 남긴다. eval_detail_docs = [] for er in eval_results: eval_detail_docs.append( { "evaluation_id": eval_id, "sample_idx": er["sample_idx"], "category_scores": er.get("category_scores", {}), "error_counts": { k: v.get("error_count", 0) for k, v in er.get("details", {}).items() if isinstance(v, dict) }, "severity": er.get("severity", ""), "score": er.get("score", 0.0), } ) for f in eval_failures: eval_detail_docs.append( { "evaluation_id": eval_id, "sample_idx": f["sample_idx"], "category_scores": {}, "error_counts": {}, "severity": None, # 0.0 이 아니라 None -- 채점 실패를 최저 점수와 구분한다. "score": None, "error": f["error"], } ) await store.insert_eval_details(eval_detail_docs) # Update run status await store.update_run( run_id, { "status": "completed", "completed_at": datetime.now(timezone.utc), }, ) return {**aggregated, "overall_score": overall_score, "failed": len(eval_failures)} aggregated = asyncio.run(run_eval()) click.echo() click.secho("Evaluation completed!", fg="green", bold=True) overall = aggregated.get("overall_score") if overall is None: # 채점이 한 건도 성공하지 못했다 -- 숫자를 찍으면 "0점"으로 읽힌다. click.secho(" Overall Score: 측정 없음 (전 샘플 채점 실패)", fg="red", bold=True) else: click.echo(f" Overall Score: {overall:.4f}") failed = aggregated.get("failed", 0) if failed: click.secho(f" 채점 실패: {failed}건 (score: null 로 기록됨)", fg="yellow") scores = aggregated.get("scores", {}) if scores: click.echo(f" Scores: {format_scores(scores)}") click.echo(f"\nResults saved to: {run_dir}/") @runs_group.command("results") @click.argument("run_ref") @click.option("--limit", "-n", default=None, type=int, help="Limit results") @click.pass_context def show_results(ctx: click.Context, run_ref: str, limit: int | None) -> None: """Show per-sample results for a run.""" cfg = ctx.obj["config"] if cfg.is_remote: # run_ref format: "project/run_id" from solar_eval.cli.client import EvalClient client = EvalClient(cfg.remote_url, cfg.timeout) parts = run_ref.split("/") if len(parts) != 2: raise click.ClickException("Remote mode: use 'project-name/run-id' format") project_id = _resolve_project_id(client, parts[0]) results = client.get(f"/api/projects/{project_id}/runs/{parts[1]}/results") else: # run_ref format: "project/run_name" parts = run_ref.split("/") if len(parts) != 2: raise click.ClickException("Local mode: use 'project-name/run-name' format") results_file = cfg.artifacts_dir(parts[0]) / parts[1] / "results.jsonl" if not results_file.exists(): raise click.ClickException(f"Results not found: {results_file}") results = [ json.loads(line) for line in results_file.read_text().strip().split("\n") if line ] if not results: click.secho("No results found.", fg="yellow") return if limit: results = results[:limit] click.secho(format_results_table(results), fg="blue") click.echo(f"\nShowing {len(results)} result(s)")