import asyncio import logging import sys from datetime import datetime, timezone from typing import Any, Callable # Ensure runner logs are visible in CLI logging.basicConfig(stream=sys.stderr, level=logging.INFO, format="%(message)s") from solar_eval.core.dataset_loader import DatasetLoader from solar_eval.evaluators.registry import create_evaluator from solar_eval.models.enums import RunStatus from solar_eval.models.sample import EvalSample from solar_eval.pipelines.registry import create_pipeline from solar_eval.providers.base import BaseProvider from solar_eval.stores.base import ResultStore logger = logging.getLogger(__name__) OnProgress = Callable[[dict[str, Any]], Any] def build_eval_sample_from_result(result: dict[str, Any], reference: Any) -> EvalSample: """저장된 결과 dict(레거시 8+trace 키)에서 채점용 `EvalSample` 을 되살린다. `insert_run_result` 가 쓰는 키 이름(`golden`/`trace`)은 `EvalSample` 필드 이름 (`reference`/`artifacts`)과 다르다 -- eval-store 하위 호환을 위해 저장 스키마는 그대로 두고(마이그레이션 계획 §9-1) 채점 직전에만 여기서 되돌린다. `execute_run` 내부 루프와 `runs eval` CLI(재채점, 새 추론 없이 디스크의 `results.jsonl` 을 그대로 씀) 양쪽에서 공유한다. Args: result: `insert_run_result` 에 넘긴 것과 같은 형태(또는 `JsonlStore. load_existing_results()` 로 디스크에서 다시 읽은 것). `input`/`output`/ `trace` 키를 읽는다. reference: 정답(golden) 값. 호출자가 직접 넘긴다 -- 두 호출자의 golden 추출 경로가 다르기 때문이다(아래 참고). 이 함수는 어느 쪽인지 모른 채 값만 받아 그대로 `sample.reference` 에 채운다. - `execute_run` 내부 루프: `_golden_raw`(디스크에 안 남는 실행 중 임시 키, resume 시에도 매 샘플마다 새로 계산됨)를 넘긴다. - `runs eval` CLI: 디스크에 저장된 `golden` 키를 그대로 넘긴다 (`_golden_raw` 는 애초에 저장되지 않아 CLI 쪽엔 없다). Returns: `input`/`output`/`reference`/`artifacts` 가 채워진 `EvalSample`. `contexts` 는 아직 아무도 안 써서 기본값(`None`) 그대로다. """ trace = result.get("trace") or {} return EvalSample( input=result.get("input"), output=result.get("output"), reference=reference, artifacts=trace, ) class BatchRunner: """Runs inference + evaluation batches with progress tracking.""" def __init__( self, store: ResultStore, inference_provider: BaseProvider, judge_provider: BaseProvider | None = None, dataset_loader: DatasetLoader | None = None, ) -> None: self.store = store self.inference_provider = inference_provider self.judge_provider = judge_provider self.dataset_loader = dataset_loader or DatasetLoader() async def execute_run( self, run_id: str, project_config: dict[str, Any], task_config: dict[str, Any], prompt: dict[str, Any], on_progress: OnProgress | None = None, max_workers: int = 5, limit: int | None = None, completed_results: list[dict[str, Any]] | None = None, ) -> None: """Execute a full inference + evaluation run. Args: completed_results: Previously completed results for resume. Samples with matching sample_idx will be skipped. 샘플 단위 실패(추론/평가)는 삼키고 상태(PARTIAL/eval_failed_count)로 기록하지만, run 을 통째로 못 돌게 만드는 예외(데이터셋 로드 실패, 파이프라인/평가기 생성 실패, 스토어 쓰기 실패 등)는 상태를 FAILED 로 남긴 뒤 그대로 재전파한다 -- "무슨 일이 있어도 예외 없이 반환"이 아니다. 호출자는 이 함수가 raise 할 수 있다고 가정해야 한다. """ try: # Build set of already-completed sample indices for resume completed_by_idx: dict[int, dict[str, Any]] = {} if completed_results: for r in completed_results: completed_by_idx[r["sample_idx"]] = r # Update status to running await self.store.update_run( run_id, { "status": RunStatus.RUNNING, "started_at": datetime.now(timezone.utc), }, ) # Load dataset dataset_config = project_config.get("dataset", {}) source = dataset_config.get("source", "huggingface") repo_name = dataset_config.get("repo", "") dataset_path = task_config.get("dataset_path", "") data = self.dataset_loader.load_jsonl(repo_name, dataset_path, source=source) if limit and limit < len(data): data = data[:limit] remaining = len(data) - len(completed_by_idx) await self.store.update_run(run_id, {"total_samples": len(data)}) if on_progress: await on_progress({"type": "dataset_loaded", "total": len(data)}) if completed_by_idx: logger.info(f"Resuming: {len(completed_by_idx)} done, {remaining} remaining") # Create pipeline # project_dir: v24 의 tool_calling_judge 가 pmi_lookup 상대경로를 풀 때만 # 쓴다 (§5-F). dataset_loader.base_dir 이 projects 루트이므로 project 이름을 # 붙이면 project_dir 이 된다 -- CLI(`_start_local`)가 이미 쓰는 것과 같은 관례. project_name = project_config.get("name") project_dir = ( self.dataset_loader.base_dir / project_name if self.dataset_loader.base_dir and project_name else None ) # config_dir: 레포 관리 프로젝트면 project_loader 가 채워 둔 config 정본 # 경로(03-evaluation). 치환 사전 같은 지식 자산이 여기서 온다. pipeline = create_pipeline( pipeline_type=task_config.get("pipeline", "single_step"), input_fields=task_config.get("input_fields", []), pipeline_config=task_config.get("pipeline_config"), prompts=prompt.get("step_prompts", {}), dataset_loader=self.dataset_loader, project_dir=project_dir, config_dir=project_config.get("config_dir"), ) # Run inference with concurrency control semaphore = asyncio.Semaphore(max_workers) completed = len(completed_by_idx) async def process_sample(idx: int, sample: dict) -> dict[str, Any]: nonlocal completed # Skip already-completed samples (resume) if idx in completed_by_idx: existing = completed_by_idx[idx] golden_field = task_config.get("golden_field") golden = sample.get(golden_field, "") if golden_field else "" return {**existing, "_golden_raw": golden} async with semaphore: input_data = { field: sample.get(field, "") for field in task_config.get("input_fields", []) } # Get golden reference golden_field = task_config.get("golden_field") golden_fields = task_config.get("golden_fields") if golden_field: golden = sample.get(golden_field, "") elif golden_fields: golden = {k: sample.get(v, "") for k, v in golden_fields.items()} else: golden = "" eval_sample = EvalSample(input=input_data, reference=golden) try: eval_sample = await pipeline.run( sample=eval_sample, prompts=prompt.get("system_prompt", ""), provider=self.inference_provider, model=prompt.get("model", "solar-pro2"), temperature=prompt.get("temperature", 0.0), max_tokens=prompt.get("max_tokens", 8000), reasoning_effort=prompt.get("reasoning_effort"), messages=prompt.get("messages"), ) except Exception as e: ts = datetime.now().strftime("%H:%M:%S") logger.warning(f"[{ts}] Sample {idx} failed: {type(e).__name__}: {e}") raise # 저장 형태는 이전과 같은 8키 dict (eval-store 의 resultRowSchema 가 # 읽는 이름들과 하위 호환) + trace 를 추가한다. trace 는 artifacts # 전체를 그대로 넣는다 -- step_outputs 뿐 아니라 v24 가 채우는 # judge_decisions/judge_tool_calls/self_consistency_runs/corrections # 도 여기 안 넣으면 저장 직전에 통째로 버려진다 (실제 judge LLM 호출· # self-consistency 반복 호출 비용이 나간 산출물이다). 특정 키만 # 하드코딩해 옮기면 다음에 파이프라인이 새 artifacts 키를 추가할 # 때마다 여기를 또 고쳐야 하므로 통째로 넘긴다 -- # resultRowSchema.trace 는 .passthrough() 라 여분 키를 그대로 받는다 # (source-schemas.ts:223-228). run_result = { "run_id": run_id, "sample_idx": idx, "input": input_data, "output": eval_sample.output, "golden": golden, "input_tokens": eval_sample.artifacts.get("usage", {}).get( "prompt_tokens", 0 ), "output_tokens": eval_sample.artifacts.get("usage", {}).get( "completion_tokens", 0 ), "inference_time_ms": eval_sample.artifacts.get("inference_time_ms", 0), "trace": { "step_outputs": {}, **eval_sample.artifacts, }, } await self.store.insert_run_result(run_result) completed += 1 ts = datetime.now().strftime("%H:%M:%S") logger.info(f"[{ts}] Sample {idx} completed ({completed}/{len(data)})") await self.store.update_run(run_id, {"completed_samples": completed}) if on_progress: await on_progress( { "type": "inference_progress", "completed": completed, "total": len(data), "sample_idx": idx, } ) return {**run_result, "_golden_raw": golden} tasks = [process_sample(i, sample) for i, sample in enumerate(data)] results = await asyncio.gather(*tasks, return_exceptions=True) # Filter out exceptions valid_results = [r for r in results if isinstance(r, dict)] errors = [(i, r) for i, r in enumerate(results) if isinstance(r, Exception)] if errors: logger.warning(f"Run {run_id}: {len(errors)}/{len(results)} samples failed") for sample_idx, err in errors: logger.warning(f" Sample {sample_idx} error: {type(err).__name__}: {err}") if on_progress: await on_progress({"type": "inference_complete", "total": len(valid_results)}) # Run evaluation await self.store.update_run(run_id, {"status": RunStatus.EVALUATING}) if on_progress: await on_progress({"type": "evaluation_start", "total": len(valid_results)}) evaluator = create_evaluator(task_config.get("evaluator", {"type": "llm_judge"})) # eval_results: 성공한 evaluate() 반환값 + "sample_idx"(정본, §0.5). # eval_failures: 채점 자체가 안 된 샘플 -- aggregate() 에는 절대 안 넘긴다 # (lcs_diff.aggregate() 처럼 r["details"]["tp"] 를 직접 읽는 구현이 실패 # 항목을 만나면 KeyError 로 죽는다, §0.1). eval_results: list[dict[str, Any]] = [] eval_failures: list[dict[str, Any]] = [] for i, result in enumerate(valid_results): # F6: enumerate 위치가 아니라 result["sample_idx"] 가 정본이다 -- # 중간 샘플이 추론에서 실패하면 valid_results 의 리스트 위치와 원본 # sample_idx 가 어긋난다. sample_idx = result["sample_idx"] try: eval_sample = build_eval_sample_from_result( result, reference=result["_golden_raw"] ) evaluator.validate_required_fields(eval_sample) eval_result = await evaluator.evaluate( sample=eval_sample, provider=self.judge_provider, ) eval_results.append({**eval_result, "sample_idx": sample_idx}) if on_progress and i % 10 == 0: await on_progress( { "type": "evaluation_progress", "completed": i + 1, "total": len(valid_results), } ) except Exception as e: logger.warning(f"Evaluation failed for sample {sample_idx}: {e}") eval_failures.append({"sample_idx": sample_idx, "error": str(e)}) # Aggregate and save evaluation -- eval_results 에는 실패 항목이 안 섞여 # 있으므로 기존 aggregate() 구현이 그대로 동작한다. aggregated = evaluator.aggregate(eval_results) # 채점에 성공한 샘플이 하나도 없으면 점수 자리를 비운다 -- aggregate() 는 빈 # 입력에 0.0 을 돌려주는데, 그건 "0점을 받았다"는 측정값이라 "측정 자체가 # 없었다"와 구분되지 않는다. judge 가 통째로 죽은 run 이 evalhub 차트에서 # 품질 급락으로 보이면 F1 을 반만 고친 셈이다. overall_score = aggregated.get("overall_score", 0.0) if eval_results else None eval_id = await self.store.create_evaluation( { "run_id": run_id, "scores": aggregated.get("scores", {}), "overall_score": overall_score, "eval_model": "gpt-4o", "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 -- 성공/실패 둘 다 한 행씩 남긴다. # results.jsonl 과 eval_details.jsonl 이 항상 같은 sample_idx 집합을 # 가리키게 해서 부분적으로만 채점된 run 에서 "이 샘플은 왜 안 보이지"를 # 없앤다. 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 -- 채점 실패를 최저 점수와 구분한다 (F1). "score": None, "error": f["error"], } ) await self.store.insert_eval_details(eval_detail_docs) # Mark run as completed/partial -- 추론이 일부 샘플에서 실패했으면 # PARTIAL, 평가 실패는 이 상태에 영향을 주지 않는다(평가 성공/실패는 # 위 eval_success_count/eval_failed_count 로만 표현한다, §0.4). errors # 의 인덱스는 asyncio.gather 가 tasks 순서를 보존하므로 이미 진짜 # sample_idx 다. final_status = RunStatus.PARTIAL if errors else RunStatus.COMPLETED await self.store.update_run( run_id, { "status": final_status, "completed_at": datetime.now(timezone.utc), "failed_samples": len(errors), "failed_sample_indices": [i for i, _ in errors], }, ) if on_progress: await on_progress( { "type": "done", "overall_score": aggregated.get("overall_score", 0.0), "scores": aggregated.get("scores", {}), } ) except Exception as e: logger.exception(f"Run {run_id} failed") await self.store.update_run( run_id, { "status": RunStatus.FAILED, "completed_at": datetime.now(timezone.utc), }, ) if on_progress: await on_progress({"type": "error", "error": str(e)}) # F4: 상태를 FAILED 로 남기고 진행 상황까지 알린 뒤 **그대로 재전파한다**. # 여기서 삼키면 호출자(CLI)는 "Results saved to ..." 를 찍고 evalhub 적재까지 # 시도한 다음 exit 0 을 낸다 -- 데이터셋조차 못 읽은 run 이 성공으로 보인다. raise