diff --git a/examples/config/starter.yaml b/examples/config/starter.yaml index ae7e04e..b45ac0b 100644 --- a/examples/config/starter.yaml +++ b/examples/config/starter.yaml @@ -149,7 +149,7 @@ default_states: trainer: train_framework: "llamafactory" - llamafactory_dir: "/home/lpc/repos/LLaMA-Factory/ + llamafactory_dir: "/home/lpc/repos/LLaMA-Factory" llamafactory_env_path: "/home/lpc/miniconda3/envs/lmf/bin/" CUDA_VISIBLE_DEVICES: "0,1" swanlab_api_key: "" diff --git a/examples/scripts/run_judger_standalone.py b/examples/scripts/run_judger_standalone.py index 448cd16..1973a8c 100644 --- a/examples/scripts/run_judger_standalone.py +++ b/examples/scripts/run_judger_standalone.py @@ -134,27 +134,13 @@ def _load_config(config_path: str) -> Dict[str, Any]: def _extract_state_from_starter_yaml(config: Dict[str, Any]) -> Dict[str, Any]: """从 starter.yaml 格式提取 state 字典。 - starter.yaml 结构: - default_states: - task_id: "..." - output_dir: "..." - judger: - eval_model_path: "..." - ... - system: - CUDA_VISIBLE_DEVICES: "..." + 仅提取 task_id 和 output_dir,不提取 judger 配置。 + Judger 配置优先从 DB (taskmodel) 读取。 """ defaults = config.get("default_states", {}) - system = config.get("system", {}) - - judger = dict(defaults.get("judger", {})) - - # 从 system 补充 GPU 配置 - if "cuda_visible_devices" not in judger and "CUDA_VISIBLE_DEVICES" in system: - judger["cuda_visible_devices"] = str(system["CUDA_VISIBLE_DEVICES"]) return { - "judger": judger, + "judger": {}, "task_id": defaults.get("task_id", ""), "output_dir": defaults.get("output_dir", "./outputs"), } diff --git a/loopai/schema/states.py b/loopai/schema/states.py index b428805..40d873d 100644 --- a/loopai/schema/states.py +++ b/loopai/schema/states.py @@ -997,14 +997,7 @@ class JudgerState(BaseModel): description="评估模型路径", json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} ) - eval_task_type: str = Field( - default="code", - title="评估任务类型", - description="评估任务类型, 支持代码生成(code), Text2sql(text2sql), 通用领域文本评估(general_text)", - json_schema_extra={"ui_type": "list", "ui_group": "评估模型", - "allowed_values": ["code", "text2sql", "general_text"]} - ) - # eval_base_url: str = Field( + #eval_base_url: str = Field( # default=None, # title="评估模型 Base URL", # description="评估模型 Base URL,未设置或为空的时候,将会尝试通过本地开启vllm", @@ -1028,13 +1021,7 @@ class JudgerState(BaseModel): description="评估模型 Top P", json_schema_extra={"ui_type": "slider", "max": 1, "ui_group": "评估模型"} ) - eval_problem_path: str = Field( - default=None, - title="评估模型问题路径", - description="评估模型问题路径", - json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} - ) - # eval_format_type: str = Field( + #eval_format_type: str = Field( # default=None, # title="评估模型问题格式化类型", # description="评估模型问题格式化类型,如果为空或None将不进入格式化节点,改格式化方式可以用户自由定义,目前支持\"human-eval\"和\"mbpp\",格式化后的文件将存至output_dir定义的目录下", @@ -1053,25 +1040,6 @@ class JudgerState(BaseModel): description="评估模型每个问题的样例生成数量", json_schema_extra={"ui_type": "number", "ui_group": "评估模型"} ) - eval_text2sql_dir: str = Field( - default=None, - title="评估模型text2sql数据库目录", - description="评估模型text2sql数据库目录,仅text2sql任务下生效,并且数据文件中需要以字段db_id标注出相应的数据库文件夹至路径目录下", - json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} - ) - # 统一vllm配置删除 - # eval_env_configs: str = Field( - # default='{"NCCL_P2P_DISABLE": "1","NCCL_IB_DISABLE": "1","NCCL_DEBUG": "INFO","NCCL_SOCKET_IFNAME": "lo","NCCL_BLOCKING_WAIT": "1"}', - # title="评估模型vllm启动环境参数", - # description="评估模型vllm启动环境参数,需要完整字符串配置,为空则认为已启动vllm将会跳过启动vllm的过程", - # json_schema_extra={"ui_type": "textarea", "language": "json", "ui_group": "评估模型"} - # ) - # eval_vllm_port: int = Field( - # default=8911, - # title="vllm本地启动参数——port", - # description="vllm本地启动参数——port,用于本地启动vllm服务的参数之一,当参数eval_base_url未设置或为空时生效", - # json_schema_extra={"ui_type": "number", "ui_group": "评估模型"} - # ) eval_vllm_tensor_parallel_size: int = Field( default=2, title="vllm本地启动参数——tensor_parallel_size", @@ -1091,41 +1059,29 @@ class JudgerState(BaseModel): # description="vllm本地启动参数——启动环境,用于本地启动vllm服务的参数之一,当参数eval_base_url未设置或为空时生效,为空时默认为当前环境启动。参数需要具体到python目录,格式应为/miniconda3/envs//bin/python", # json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} # ) - output_result_path: str = Field( - default="", - title="评测结果文件保存路径", - description="评测结果文件保存路径,该参数不支持用户自定义,运行后由程序根据任务ID等参数生成", - json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} - ) - output_case_path: str = Field( - default="", - title="评测样例集文件保存路径", - description="评测样例集文件保存路径,该参数不支持用户自定义,运行后由程序根据任务ID等参数生成", - json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} - ) - output_problem_path: str = Field( - default="", - title="评测格式化后问题集保存路径", - description="评测格式化后问题集,该参数不支持用户自定义,运行后由程序根据任务ID等参数生成,如未使用格式化模版该路径即为原始问题文件的路径", - json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} + benchlist: List[Dict[str, Any]] = Field( + default_factory=list, + title="主任务评测集", + description="主任务评测集列表,每个元素包含 name、task_type、problem_path 等字段", + json_schema_extra={"ui_type": "textarea", "ui_group": "评估模型"} ) - output_pred_path: str = Field( - default="", - title="评测预测结果保存路径", - description="通用文本评测结束后产生的预测文件路径", - json_schema_extra={"ui_type": "file_path", "ui_group": "评估模型"} + extra_benchlist: List[Dict[str, Any]] = Field( + default_factory=list, + title="附加任务评测集", + description="附加任务评测集列表,格式同 benchlist。失败不影响主任务", + json_schema_extra={"ui_type": "textarea", "ui_group": "评估模型"} ) - bench: Dict[str, Any] = Field( - default_factory=dict, - title="Bench运行信息", - description="通用文本评测生成的 bench 信息,用于 Analyzer 后续指标计算", - json_schema_extra={"ui_type": "json_viewer", "ui_group": "评估模型"} + bench_result: List[Dict[str, Any]] = Field( + default_factory=list, + title="主任务评测结果", + description="主任务 bench 评测结果列表,供 Analyzer 读取", + json_schema_extra={"ui_type": "textarea", "ui_group": "评估模型"} ) - bench_name: str = Field( - default="general_text_eval", - title="评测集名称", - description="通用文本评测使用的评测集名称", - json_schema_extra={"ui_type": "text", "ui_group": "评估模型"} + extra_bench_result: List[Dict[str, Any]] = Field( + default_factory=list, + title="附加任务评测结果", + description="附加任务 bench 评测结果列表,供 Analyzer 读取", + json_schema_extra={"ui_type": "textarea", "ui_group": "评估模型"} ) cuda_visible_devices: str = Field( default="0", @@ -1141,37 +1097,7 @@ class JudgerState(BaseModel): # description="是否通过 API 调用模型", # json_schema_extra={"ui_type": "toggle_switch", "ui_group": "评估模型"} # ) - bench: List[Dict[str, Any]] = Field( - default="", - title="评测集名称", - description="通用文本评测使用的评测集相关信息", - json_schema_extra={"ui_type": "textarea", "ui_group": "评估模型"} - ) - bench_dataflow_eval_type: str = Field( - default="", - title="通用文本评测类型", - description="通用文本 One-Eval DataFlow 评测类型,例如 key2_qa / key1_text_score", - json_schema_extra={"ui_type": "list", "ui_group": "评估模型", - "allowed_values": ["key1_text_score", "key2_qa", "key2_q_ma", "key3_q_choices_a", "key3_q_choices_as", "key3_q_a_rejected"]} - ) - key_mapping: Dict[str, Any] = Field( - default_factory=dict, - title="字段映射", - description="DataFlow 评测字段映射,如 input_question_key / input_target_key / input_pred_key", - json_schema_extra={"ui_type": "json_viewer", "ui_group": "评估模型"} - ) - # skip_dataflow_eval: bool = Field( - # default=False, - # title="跳过 DataFlow 正式评测", - # description="为 True 时仅准备 bench / records,不调用 DataFlowEvalTool.run_eval", - # json_schema_extra={"ui_type": "toggle_switch", "ui_group": "评估模型"} - # ) - # output_dir: str = Field( - # default="", - # title="通用文本输出路径", - # description="通用文本任务结束后输出路径", - # json_schema_ectra={"ui_type": "text", "ui_group": "评估模型"} - # ) + class AnalyzerState(BaseModel): diff --git a/loopai/skills/Judger/__init__.py b/loopai/skills/Judger/__init__.py index 99f9cf3..8f9a5ef 100644 --- a/loopai/skills/Judger/__init__.py +++ b/loopai/skills/Judger/__init__.py @@ -64,20 +64,27 @@ def run( # 成功——标准 payload 输出到 stdout(Codex 消费) judger = result.get("judger", {}) - bench = judger.get("bench") or {} - metrics = judger.get("metrics") or {} - if not metrics: - metrics = (bench.get("meta") or {}).get("eval_result") or {} + bench_list = judger.get("bench_result") or [] + secondary_list = judger.get("extra_bench_result") or [] + + # 聚合 metrics:按 bench_name 索引 + metrics = {} + for b in bench_list + secondary_list: + m = b.get("metrics") or {} + if m: + metrics[b.get("bench_name", "unknown")] = m + else: + meta = b.get("meta") or {} + m = meta.get("eval_result") or {} + if m: + metrics[b.get("bench_name", "unknown")] = m + metrics_str = json.dumps(metrics, ensure_ascii=False) if metrics else "" emit_success( data={ - "task_type": judger.get("eval_task_type"), - "output_result_path": judger.get("output_result_path", ""), - "output_case_path": judger.get("output_case_path", ""), - "output_problem_path": judger.get("output_problem_path", ""), - "output_pred_path": judger.get("output_pred_path", ""), - "bench": bench, + "bench_result": bench_list, + "extra_bench_result": secondary_list, "metrics": metrics_str, }, stream_writer=writer, diff --git a/loopai/skills/Judger/cli.py b/loopai/skills/Judger/cli.py new file mode 100644 index 0000000..47525cd --- /dev/null +++ b/loopai/skills/Judger/cli.py @@ -0,0 +1,29 @@ +# -*- coding: utf-8 -*- +"""Judger CLI entry point — ``loopai-judger`` command.""" + +from __future__ import annotations + +import argparse + + +def main(): + parser = argparse.ArgumentParser( + description="Run LoopAI Judger evaluation pipeline (standalone, no LangGraph)", + ) + parser.add_argument( + "--resume", action="store_true", default=False, + help="Resume from last checkpoint", + ) + parser.add_argument( + "--from-step", type=str, default=None, + help="Force start from a specific pipeline step", + ) + + args = parser.parse_args() + + from loopai.skills.Judger import run + run(resume=args.resume, from_step=args.from_step) + + +if __name__ == "__main__": + main() diff --git a/loopai/skills/Judger/runner.py b/loopai/skills/Judger/runner.py index b06417b..3168709 100644 --- a/loopai/skills/Judger/runner.py +++ b/loopai/skills/Judger/runner.py @@ -161,8 +161,7 @@ def _save_task_progress(state: Dict[str, Any], task_id: str) -> None: judger = state.get("judger", {}) updates: Dict[str, Any] = {} - for k in ("output_result_path", "output_case_path", "output_problem_path", - "output_pred_path", "bench"): + for k in ("bench_result", "extra_bench_result"): if k in judger and judger[k]: updates[k] = judger[k] update_configer_task_state_config("judger", updates, task_id=task_id) @@ -282,6 +281,8 @@ def _step_validate(state: Dict[str, Any], writer) -> Dict[str, Any]: ) # 5. JSONL 字段校验 + from loopai.skills.Judger.utils.data import check_jsonl_fields + if task_type == "code": fmt = judger.get("eval_format_type", "") if fmt == "mbpp": @@ -309,11 +310,6 @@ def _step_validate(state: Dict[str, Any], writer) -> Dict[str, Any]: message=f"Problem file {problem_path} has invalid fields for task type {task_type}.", ) - # 6. 重置输出路径 - state["judger"]["output_result_path"] = "" - state["judger"]["output_case_path"] = "" - state["judger"]["output_problem_path"] = "" - writer(StreamEvent( current=state.get("current"), progress=1.0, message="配置校验通过", data={"task_type": task_type, "problem_path": problem_path})) @@ -507,6 +503,92 @@ def _run_step(step_name: str, state: Dict[str, Any], writer) -> Dict[str, Any]: ) +def _apply_bench_to_state(state: Dict[str, Any], bench: Dict[str, Any]) -> None: + """将 bench entry 的字段注入到 state["judger"],使标准 pipeline 可直接运行。""" + judger = state.setdefault("judger", {}) + # 清除上一个 bench 的字段,避免残留 + for k in ("eval_format_type", "eval_text2sql_dir", + "bench_dataflow_eval_type", "key_mapping"): + judger.pop(k, None) + judger["eval_task_type"] = bench.get("task_type", "code") + judger["eval_problem_path"] = bench.get("problem_path", "") + judger["bench_name"] = bench.get("name", "") + if bench.get("case_num") is not None: + judger["eval_case_num"] = bench["case_num"] + else: + judger.setdefault("eval_case_num", 10) + if bench.get("batch_size") is not None: + judger["eval_batch_size"] = bench["batch_size"] + else: + judger.setdefault("eval_batch_size", 10) + if bench.get("text2sql_dir"): + judger["eval_text2sql_dir"] = bench["text2sql_dir"] + if bench.get("eval_type"): + judger["bench_dataflow_eval_type"] = bench["eval_type"] + if bench.get("key_mapping"): + judger["key_mapping"] = bench["key_mapping"] + + +def _run_single_bench( + state: Dict[str, Any], + bench: Dict[str, Any], + writer, +) -> Dict[str, Any]: + """运行单个 bench 的完整流水线,返回 bench result dict。""" + _apply_bench_to_state(state, bench) + + task_type = bench["task_type"] + if task_type == "general_text": + steps = _GENERAL_TEXT_STEPS + else: + steps = _CODE_TEXTSQL_STEPS + + bench_name = bench["name"] + logger.info(f"[Judger] bench {bench_name} (task_type={task_type}) starting...") + writer(StreamEvent( + current="judger", progress=0.0, + message=f"Bench 开始: {bench_name}", + data={"bench_name": bench_name, "task_type": task_type})) + + for step_name in steps: + state["current"] = f"{bench_name}.{step_name}" + logger.info(f"[Judger] [{bench_name}] step {step_name}") + if step_name == "finish": + break + state = _run_step(step_name, state, writer) + + # 收集结果 + judger = state.get("judger", {}) + if task_type == "general_text": + bench_data = judger.get("bench") or {} + result = { + "bench_name": bench_name, + "task_type": task_type, + "output_result_path": judger.get("output_result_path", ""), + "output_pred_path": judger.get("output_pred_path", ""), + "eval_status": bench_data.get("eval_status", "success"), + "meta": bench_data.get("meta", {}), + "key_mapping": bench_data.get("key_mapping", {}), + "metrics": (bench_data.get("meta", {})).get("eval_result", {}), + } + else: + result = { + "bench_name": bench_name, + "task_type": task_type, + "output_case_path": judger.get("output_case_path", ""), + "output_result_path": judger.get("output_result_path", ""), + "metrics": judger.get("metrics", {}), + "eval_status": "success", + } + + writer(StreamEvent( + current="judger", progress=1.0, + message=f"Bench 完成: {bench_name}", + data={"bench_name": bench_name, "result": result})) + logger.info(f"[Judger] bench {bench_name} done") + return result + + # --------------------------------------------------------------------------- # Main pipeline runner / 主流水线执行器 # --------------------------------------------------------------------------- @@ -543,6 +625,20 @@ def run_judger_pipeline( state = _load_task_state(task_id) elif state is not None: state = dict(state) + # If state from starter yaml is incomplete (e.g. missing benchlist), + # load full judger config from DB (task model state) + try: + db_state = _load_task_state(task_id) + db_judger = db_state.get("judger") or {} + if db_judger: + # Merge DB judger fields into state, but keep output_dir from yaml + db_judger.pop("_last_completed", None) + db_judger.pop("_current", None) + for k, v in db_judger.items(): + if v is not None and v != "": + state["judger"][k] = v + except Exception: + pass else: state = _load_task_state(task_id) if not state.get("judger"): @@ -571,63 +667,101 @@ def run_judger_pipeline( output_dir = state.get("output_dir", "./outputs") - writer(StreamEvent( - current="judger", progress=0.0, message="Judger pipeline started", - data={"task_id": task_id, "resume": resume, - "task_type": (state.get("judger") or {}).get("eval_task_type")})) - - # 根据任务类型选择流水线 - task_type = (state.get("judger") or {}).get("eval_task_type", "code") - steps = _GENERAL_TEXT_STEPS if task_type == "general_text" else _CODE_TEXTSQL_STEPS - judger_cfg = state.get("judger") or {} - logger.info(f"[Judger] task_id={task_id} task_type={task_type} " - f"pipeline={' -> '.join(steps)}") - logger.info(f"[Judger] model_path={judger_cfg.get('eval_model_path')} " - f"problem_path={judger_cfg.get('eval_problem_path')} " - f"batch_size={judger_cfg.get('eval_batch_size')} " - f"case_num={judger_cfg.get('eval_case_num')} " - f"gpu={judger_cfg.get('cuda_visible_devices')} " - f"tp_size={judger_cfg.get('eval_vllm_tensor_parallel_size')}") - - # 确定起始步骤 - if from_step is not None: - start_step = normalize_judger_step(from_step) - elif resume: - start_step = _resume_step_from_state(state) - else: - start_step = steps[0] - if start_step is None: - start_step = steps[0] - start_at = _start_index(start_step, steps) + def _parse_benchlist(val): + """textarea 可能返回 JSON 字符串(单行或多行),转为 list。""" + if isinstance(val, str): + # 尝试整体解析 + try: + parsed = json.loads(val) + if isinstance(parsed, list): + return parsed + except (json.JSONDecodeError, ValueError): + pass + # 尝试按行解析(每行一个 JSON 对象) + items = [] + for line in val.strip().split("\n"): + line = line.strip() + if line: + try: + items.append(json.loads(line)) + except (json.JSONDecodeError, ValueError): + pass + if items: + return items + return [] + return val if isinstance(val, list) else [] + + benchlist = _parse_benchlist(judger_cfg.get("benchlist")) or [] + extra_benchlist = _parse_benchlist(judger_cfg.get("extra_benchlist")) or [] + + if not benchlist and not extra_benchlist: + emit_error( + ValueError("benchlist 和 extra_benchlist 都为空,请至少配置一个评测集"), + code=ErrorCode.CONFIG_ERROR, recoverable=True, + stream_writer=writer, + message="Both benchlist and extra_benchlist are empty. Please configure at least one bench.", + ) - if from_step is None and _is_finished(state): - logger.info(f"[Judger] already finished, skip") - return state + writer(StreamEvent( + current="judger", progress=0.0, message="Judger pipeline started", + data={"task_id": task_id, "resume": resume})) - # 按序执行各步骤 - for i, step_name in enumerate(steps[start_at:], start_at): - state["current"] = f"judger.{step_name}" - logger.info(f"[Judger] step [{i+1}/{len(steps)}] {step_name} starting...") - _save_task_progress(state, task_id) + logger.info(f"[Judger] task_id={task_id} " + f"benchlist={[b.get('name') for b in benchlist]} " + f"extra_benchlist={[b.get('name') for b in extra_benchlist]}") - writer(StreamEvent( - current=state["current"], progress=0.0, - message=f"步骤开始: {step_name}")) + bench_results: List[Dict[str, Any]] = [] + secondary_results: List[Dict[str, Any]] = [] + state["judger"]["bench_result"] = bench_results + state["judger"]["extra_bench_result"] = secondary_results - if step_name == "finish": - state["last_completed"] = "finish" + # 主任务(失败记录后退出) + for bench in benchlist: + try: + result = _run_single_bench(state, bench, writer) + bench_results.append(result) _save_task_progress(state, task_id) - writer(StreamEvent( - current=state["current"], progress=1.0, message="流水线完成")) - logger.info(f"[Judger] pipeline finished") - return state - - state = _run_step(step_name, state, writer) - state["last_completed"] = step_name - _save_task_progress(state, task_id) - - logger.info(f"[Judger] step [{i+1}/{len(steps)}] {step_name} done") + except SystemExit: + bench_results.append({ + "bench_name": bench.get("name", "unknown"), + "eval_status": "failed", + "meta": {"error": "Bench evaluation failed"}, + }) + _save_task_progress(state, task_id) + raise + except Exception as exc: + bench_results.append({ + "bench_name": bench.get("name", "unknown"), + "eval_status": "failed", + "meta": {"error": str(exc)}, + }) + _save_task_progress(state, task_id) + raise + # 附加任务(失败记录后继续) + for bench in extra_benchlist: + try: + result = _run_single_bench(state, bench, writer) + secondary_results.append(result) + _save_task_progress(state, task_id) + except SystemExit: + secondary_results.append({ + "bench_name": bench.get("name", "unknown"), + "eval_status": "failed", + "meta": {"error": "Bench evaluation failed"}, + }) + except Exception as exc: + secondary_results.append({ + "bench_name": bench.get("name", "unknown"), + "eval_status": "failed", + "meta": {"error": str(exc)}, + }) + + state["last_completed"] = "finish" + _save_task_progress(state, task_id) + writer(StreamEvent( + current="finish", progress=1.0, message="流水线完成")) + logger.info(f"[Judger] pipeline finished") return state diff --git a/loopai/skills/Judger/runtime_config.py b/loopai/skills/Judger/runtime_config.py index 72922e6..2f37953 100644 --- a/loopai/skills/Judger/runtime_config.py +++ b/loopai/skills/Judger/runtime_config.py @@ -5,9 +5,7 @@ _DEFAULT_OUTPUT_DIR = "./outputs" -# JudgerState schema defaults (mirror loopai/schema/states.py JudgerState) _SCHEMA_DEFAULTS: Dict[str, Any] = { - "eval_task_type": "code", "eval_temperature": 0, "eval_top_p": 0.95, "eval_batch_size": 10, @@ -15,9 +13,6 @@ "eval_vllm_tensor_parallel_size": 1, "eval_vllm_gpu_memory_utilization": 0.9, "cuda_visible_devices": "0", - "bench_name": "general_text_eval", - "bench_dataflow_eval_type": "", - "key_mapping": {}, } @@ -42,23 +37,19 @@ def resolve_judger_runtime_config( state: Optional[Dict[str, Any]], task_id: Optional[str] = None, ) -> Dict[str, Any]: - """Resolve Judger runtime values from env, state, defaults. + """Resolve Judger global runtime config from state + env + defaults. - Priority: env > state["judger"] > schema defaults. + bench-level fields (name, task_type, problem_path, eval_type, etc.) + come from ``benchlist`` / ``extra_benchlist`` directly. """ judger = _judger(state) is_state_dict = isinstance(state, dict) - # --- model / vllm --- + # --- model / vllm (global) --- model_path = _first_non_empty( os.getenv("JUDGER_MODEL_PATH"), judger.get("eval_model_path"), ) - task_type = _first_non_empty( - os.getenv("JUDGER_TASK_TYPE"), - judger.get("eval_task_type"), - _SCHEMA_DEFAULTS["eval_task_type"], - ) temperature = _first_non_empty( os.getenv("JUDGER_TEMPERATURE"), judger.get("eval_temperature"), @@ -69,10 +60,6 @@ def resolve_judger_runtime_config( judger.get("eval_top_p"), _SCHEMA_DEFAULTS["eval_top_p"], ) - problem_path = _first_non_empty( - os.getenv("JUDGER_PROBLEM_PATH"), - judger.get("eval_problem_path"), - ) batch_size = _first_non_empty( os.getenv("JUDGER_BATCH_SIZE"), judger.get("eval_batch_size"), @@ -83,12 +70,9 @@ def resolve_judger_runtime_config( judger.get("eval_case_num"), _SCHEMA_DEFAULTS["eval_case_num"], ) - - # --- vllm local startup --- tensor_parallel_size = _first_non_empty( os.getenv("JUDGER_TENSOR_PARALLEL_SIZE"), judger.get("eval_vllm_tensor_parallel_size"), - judger.get("tensor_parallel_size"), # 兼容 starter.yaml 中的短键名 _SCHEMA_DEFAULTS["eval_vllm_tensor_parallel_size"], ) gpu_memory_utilization = _first_non_empty( @@ -102,30 +86,8 @@ def resolve_judger_runtime_config( _SCHEMA_DEFAULTS["cuda_visible_devices"], ) - # --- general_text --- - bench_name = _first_non_empty( - os.getenv("JUDGER_BENCH_NAME"), - judger.get("bench_name"), - _SCHEMA_DEFAULTS["bench_name"], - ) - bench_dataflow_eval_type = _first_non_empty( - os.getenv("JUDGER_BENCH_DATAFLOW_EVAL_TYPE"), - judger.get("bench_dataflow_eval_type"), - _SCHEMA_DEFAULTS["bench_dataflow_eval_type"], - ) - key_mapping = _first_non_empty( - judger.get("key_mapping"), - _SCHEMA_DEFAULTS["key_mapping"], - ) - if isinstance(key_mapping, str): - import json - try: - key_mapping = json.loads(key_mapping) - except (json.JSONDecodeError, TypeError): - key_mapping = {} - # --- global --- - task_id = _first_non_empty( + resolved_task_id = _first_non_empty( task_id, os.getenv("TASK_ID"), state.get("task_id") if is_state_dict else None, @@ -165,27 +127,19 @@ def resolve_judger_runtime_config( # --- write resolved values back into state --- if is_state_dict: - # CLI/env 传入的 task_id 优先级高于 YAML 默认值,始终覆盖 - if task_id: - state["task_id"] = task_id + if resolved_task_id: + state["task_id"] = resolved_task_id if output_dir: state["output_dir"] = output_dir - - # 只在值非空时才覆盖,避免把 DB/state 中的已有值冲掉 for key, val in ( ("eval_model_path", model_path), - ("eval_task_type", task_type), ("eval_temperature", temperature), ("eval_top_p", top_p), - ("eval_problem_path", problem_path), ("eval_batch_size", batch_size), ("eval_case_num", case_num), ("eval_vllm_tensor_parallel_size", tensor_parallel_size), ("eval_vllm_gpu_memory_utilization", gpu_memory_utilization), ("cuda_visible_devices", cuda_visible_devices), - ("bench_name", bench_name), - ("bench_dataflow_eval_type", bench_dataflow_eval_type), - ("key_mapping", key_mapping), ): if val is not None: state["judger"][key] = val @@ -193,20 +147,15 @@ def resolve_judger_runtime_config( state["DB_PATH"] = db_path return { - "task_id": str(task_id) if task_id else "", + "task_id": str(resolved_task_id) if resolved_task_id else "", "output_dir": str(output_dir or _DEFAULT_OUTPUT_DIR), "db_path": db_path, "model_path": model_path, - "task_type": str(task_type), "temperature": temperature, "top_p": top_p, - "problem_path": problem_path, "batch_size": batch_size, "case_num": case_num, "tensor_parallel_size": tensor_parallel_size, "gpu_memory_utilization": gpu_memory_utilization, "cuda_visible_devices": str(cuda_visible_devices), - "bench_name": str(bench_name), - "bench_dataflow_eval_type": str(bench_dataflow_eval_type or ""), - "key_mapping": key_mapping if isinstance(key_mapping, dict) else {}, } diff --git a/loopai/skills/Judger/utils/eval_general_text.py b/loopai/skills/Judger/utils/eval_general_text.py index 6878ea5..8ed2957 100644 --- a/loopai/skills/Judger/utils/eval_general_text.py +++ b/loopai/skills/Judger/utils/eval_general_text.py @@ -87,8 +87,8 @@ def _build_model_config(cfg: Dict[str, Any]) -> ModelConfig: def _generate_key_mapping(cfg: Dict[str, Any]) -> Dict[str, Any]: key_mapping: Dict[str, Any] = {} - eval_type = cfg["bench_dataflow_eval_type"] - eval_problem_path = cfg["eval_problem_path"] + eval_type = cfg.get("bench_dataflow_eval_type") or cfg.get("eval_type", "") + eval_problem_path = cfg.get("eval_problem_path") or cfg.get("problem_path", "") count = 0 with open(eval_problem_path, "r", encoding="utf-8") as f: for line in f: @@ -303,33 +303,31 @@ def _check_output_health( # --------------------------------------------------------------------------- def run_eval_general_text(state: Dict[str, Any], writer) -> Dict[str, Any]: - """通用文本评测(One-Eval DataFlowEvalTool),无 LangGraph 依赖。 - - 进度事件通过 ``writer`` 发送(PickleEventWriter), - 评测结果存入 ``state["judger"]["bench"]`` 和 ``output_result_path``。 - """ + """通用文本评测(One-Eval DataFlowEvalTool),无 LangGraph 依赖。""" os.environ["CUDA_VISIBLE_DEVICES"] = (state.get("judger") or {}).get( "cuda_visible_devices", "2" ) cfg: Dict[str, Any] = state.get("judger") or {} + + bench_name = cfg.get("bench_name") or "general_text_eval" outdir = ( Path(state.get("output_dir") or "./outputs") / (state.get("task_id") or "default_task") / "judger" / (writer.version_id or "") + / bench_name ) outdir.mkdir(parents=True, exist_ok=True) run_ts = time.strftime("%Y%m%d_%H%M%S") - # 校验必填字段 eval_result_path = cfg.get("eval_problem_path") if not eval_result_path: emit_error( - ValueError("缺少评测输入路径:请提供 judger.eval_problem_path"), + ValueError("缺少评测输入路径"), code=ErrorCode.CONFIG_ERROR, recoverable=True, stream_writer=writer, - message="Missing judger.eval_problem_path for general_text evaluation.", + message="Missing eval_problem_path for general_text evaluation.", ) if not os.path.exists(eval_result_path): emit_error( @@ -342,10 +340,10 @@ def run_eval_general_text(state: Dict[str, Any], writer) -> Dict[str, Any]: eval_type = cfg.get("bench_dataflow_eval_type") if not eval_type: emit_error( - ValueError("通用文本评测缺少 eval_type。请在 judger.bench_dataflow_eval_type 中提供"), + ValueError("通用文本评测缺少 eval_type"), code=ErrorCode.CONFIG_ERROR, recoverable=True, stream_writer=writer, - message="Missing judger.bench_dataflow_eval_type for general_text evaluation.", + message="Missing bench_dataflow_eval_type for general_text evaluation.", ) logger.info(f"[Judger/general_text] model={cfg.get('eval_model_path')} " @@ -353,7 +351,6 @@ def run_eval_general_text(state: Dict[str, Any], writer) -> Dict[str, Any]: f"tp_size={cfg.get('eval_vllm_tensor_parallel_size', 1)} " f"problem_path={eval_result_path}") - # 解析 key_mapping key_mapping = cfg.get("key_mapping") or {} if isinstance(key_mapping, str): try: @@ -363,8 +360,6 @@ def run_eval_general_text(state: Dict[str, Any], writer) -> Dict[str, Any]: if not key_mapping: key_mapping = _generate_key_mapping(cfg) - bench_name = cfg.get("bench_name") or "general_text_eval" - writer(StreamEvent( current=state.get("current", "judger"), progress=0.0, message="开始通用文本评测 (One-Eval)", @@ -542,7 +537,6 @@ def run_eval_general_text(state: Dict[str, Any], writer) -> Dict[str, Any]: bench.meta["artifact_paths"]["records_path"] = step2_file_path bench.meta["eval_detail_path"] = step2_file_path - # 回写 state state["judger"]["bench"] = { "bench_name": bench.bench_name, "dataset_cache": bench.dataset_cache, diff --git a/loopai/skills/Judger/utils/evaluate.py b/loopai/skills/Judger/utils/evaluate.py index f669d6f..60ab494 100644 --- a/loopai/skills/Judger/utils/evaluate.py +++ b/loopai/skills/Judger/utils/evaluate.py @@ -128,12 +128,14 @@ def run_evaluate_code(state: Dict[str, Any], writer) -> Dict[str, Any]: judger_state = state.get("judger", {}) output_dir = Path(state.get("output_dir")) problem_path = judger_state["eval_problem_path"] - problem_file_name = Path(problem_path).stem + bench_name = judger_state.get("bench_name", Path(problem_path).stem) test_case_path = str( - output_dir / str(state_task_id) / "judger" / writer.version_id / f"{problem_file_name}_sample.jsonl" + output_dir / str(state_task_id) / "judger" / writer.version_id + / bench_name / f"{bench_name}_sample.jsonl" ) result_path = str( - output_dir / str(state_task_id) / "judger" / writer.version_id / f"{problem_file_name}_result.jsonl" + output_dir / str(state_task_id) / "judger" / writer.version_id + / bench_name / f"{bench_name}_result.jsonl" ) case_num = judger_state.get("eval_case_num", 10) task_type = judger_state["eval_task_type"] @@ -200,12 +202,14 @@ def run_evaluate_text2sql(state: Dict[str, Any], writer) -> Dict[str, Any]: judger_state = state.get("judger", {}) output_dir = Path(state.get("output_dir")) problem_path = judger_state["eval_problem_path"] - problem_file_name = Path(problem_path).stem + bench_name = judger_state.get("bench_name", Path(problem_path).stem) test_case_path = str( - output_dir / str(state_task_id) / "judger" / writer.version_id / f"{problem_file_name}_sample.jsonl" + output_dir / str(state_task_id) / "judger" / writer.version_id + / bench_name / f"{bench_name}_sample.jsonl" ) result_path = str( - output_dir / str(state_task_id) / "judger" / writer.version_id / f"{problem_file_name}_result.jsonl" + output_dir / str(state_task_id) / "judger" / writer.version_id + / bench_name / f"{bench_name}_result.jsonl" ) case_num = judger_state.get("eval_case_num", 10) task_type = judger_state["eval_task_type"] diff --git a/loopai/skills/Judger/utils/generate.py b/loopai/skills/Judger/utils/generate.py index a4f415c..cff7428 100644 --- a/loopai/skills/Judger/utils/generate.py +++ b/loopai/skills/Judger/utils/generate.py @@ -1,9 +1,5 @@ # -*- coding: utf-8 -*- -"""Standalone code/text2sql sample generation — no LangGraph dependency. - -Extracted from ``loopai.agents.Judger.utils.oj.generate``, -replaced ``get_stream_writer()`` with a passed-in ``writer`` parameter. -""" +"""Standalone code/text2sql sample generation — no LangGraph dependency.""" from pathlib import Path from typing import Any, Dict @@ -46,9 +42,10 @@ def run_generate_code(state: Dict[str, Any], writer) -> str: output_dir = Path(state.get("output_dir")) problem_path = judger_state["eval_problem_path"] - problem_file_name = Path(problem_path).stem + bench_name = judger_state.get("bench_name", Path(problem_path).stem) test_case_path = str( - output_dir / str(state_task_id) / "judger" / writer.version_id / f"{problem_file_name}_sample.jsonl" + output_dir / str(state_task_id) / "judger" / writer.version_id + / bench_name / f"{bench_name}_sample.jsonl" ) batch_size = judger_state["eval_batch_size"] @@ -63,7 +60,8 @@ def run_generate_code(state: Dict[str, Any], writer) -> str: writer(StreamEvent( current=state.get("current", "judger"), progress=0.0, message=f"{task_type}任务样本合成开始", - data={"total_tasks": total_tasks, "num_samples_per_task": num_samples_per_task, + data={"bench_name": bench_name, "total_tasks": total_tasks, + "num_samples_per_task": num_samples_per_task, "total_samples": total_samples})) logger.info(f"===== 开始生成样本 =====") @@ -91,14 +89,14 @@ def run_generate_code(state: Dict[str, Any], writer) -> str: current=state.get("current", "judger"), progress=round(cnt / total_samples, 1), message=f"{task_type}任务样本合成进度", - data={"progress_detail": f"{cnt}/{total_samples}"})) + data={"bench_name": bench_name, "progress_detail": f"{cnt}/{total_samples}"})) write_jsonl(test_case_path, samples) logger.info(f"===== 生成完成 ===== 样本数:{len(samples)} 路径:{test_case_path}") writer(StreamEvent( current=state.get("current", "judger"), progress=1.0, message=f"{task_type}任务样本合成完成", - data={"sample_num": len(samples), "sample_save_path": test_case_path})) + data={"bench_name": bench_name, "sample_num": len(samples), "sample_save_path": test_case_path})) return test_case_path @@ -117,9 +115,10 @@ def run_generate_text2sql(state: Dict[str, Any], writer) -> str: output_dir = Path(state.get("output_dir")) problem_path = judger_state["eval_problem_path"] - problem_file_name = Path(problem_path).stem + bench_name = judger_state.get("bench_name", Path(problem_path).stem) test_case_path = str( - output_dir / str(state_task_id) / "judger" / writer.version_id / f"{problem_file_name}_sample.jsonl" + output_dir / str(state_task_id) / "judger" / writer.version_id + / bench_name / f"{bench_name}_sample.jsonl" ) task_type = judger_state["eval_task_type"] @@ -135,7 +134,8 @@ def run_generate_text2sql(state: Dict[str, Any], writer) -> str: writer(StreamEvent( current=state.get("current", "judger"), progress=0.0, message=f"{task_type}任务样本合成开始", - data={"total_tasks": total_tasks, "num_samples_per_task": num_samples_per_task, + data={"bench_name": bench_name, "total_tasks": total_tasks, + "num_samples_per_task": num_samples_per_task, "total_samples": total_samples})) logger.info(f"===== 开始生成样本 =====") @@ -170,12 +170,12 @@ def run_generate_text2sql(state: Dict[str, Any], writer) -> str: current=state.get("current", "judger"), progress=round(cnt / total_samples, 1), message=f"{task_type}任务样本合成进度", - data={"progress_detail": f"{cnt}/{total_samples}"})) + data={"bench_name": bench_name, "progress_detail": f"{cnt}/{total_samples}"})) write_jsonl(test_case_path, samples) logger.info(f"===== 生成完成 ===== 样本数:{len(samples)} 路径:{test_case_path}") writer(StreamEvent( current=state.get("current", "judger"), progress=1.0, message=f"{task_type}任务样本合成完成", - data={"sample_num": len(samples), "sample_save_path": test_case_path})) + data={"bench_name": bench_name, "sample_num": len(samples), "sample_save_path": test_case_path})) return test_case_path diff --git a/loopai/skills/Judger/utils/vllm_starter.py b/loopai/skills/Judger/utils/vllm_starter.py index 72661c5..37301e2 100644 --- a/loopai/skills/Judger/utils/vllm_starter.py +++ b/loopai/skills/Judger/utils/vllm_starter.py @@ -106,7 +106,7 @@ def start_vllm_openai_api_server( stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, text=True, - timeout=10.0, + timeout=30.0, ) except FileNotFoundError: raise Exception(f"未找到 Python 可执行文件:{python_exec}") diff --git a/setup.py b/setup.py index 6380afb..9bec3dc 100644 --- a/setup.py +++ b/setup.py @@ -65,6 +65,7 @@ entry_points={ "console_scripts": [ "loopai-obtainercli=loopai.skills.ObtainerCLI.cli:main", + "loopai-judger=loopai.skills.Judger.cli:main", ], }, python_requires=">=3.12", diff --git a/skills/Judger/SKILL.md b/skills/Judger/SKILL.md index 2cb4571..90378b0 100644 --- a/skills/Judger/SKILL.md +++ b/skills/Judger/SKILL.md @@ -2,641 +2,184 @@ ## Purpose -Judger Skill 用于在无 LangGraph(独立模式)下运行 LoopAI 评测流水线。支持三种任务类型: +无 LangGraph 的独立评测流水线。支持三种任务类型: -- **code** — 代码生成评测(human-eval / mbpp 格式),计算 pass@k -- **text2sql** — SQL 生成评测,在 SQLite 数据库上执行校验 +- **code** — 代码生成评测(human-eval / mbpp),计算 pass@k +- **text2sql** — SQL 生成评测,SQLite 执行校验 - **general_text** — 通用文本评测(One-Eval DataFlowEvalTool) -评测结果写入文件系统,进度事件持久化到 pickle,state 通过 Configer 读写 TaskModel.state。 - ## How to Invoke **唯一入口:`loopai.skills.Judger.run()`** -`DB_PATH` 和 `TASK_ID` 从环境变量自动获取,无需传参。Codex 通过子进程调用: +`DB_PATH` 和 `TASK_ID` 从环境变量自动获取: ```bash DB_PATH=api/db/db.sqlite3 TASK_ID= \ python -c "from loopai.skills.Judger import run; run()" ``` -缺环境变量时 `emit_error` 退出并输出错误 JSON。 - -**禁止使用的路径:** -- ❌ `loopai.agents.Judger.JudgerAgent` — **已删除**,请使用 `loopai.skills.Judger.run()` -- ❌ Codex 进程内直接 import `loopai.skills.Judger` — pipeline 内 `emit_error` 会 `sys.exit(1)`,杀死 Codex 自身进程 - -**正确调用后的产物(用于判断是否走了正确路径):** -- `outputs//judger.pkl` — 事件流(含 `metrics`) -- `outputs//judger//` — 评测结果文件(每次运行独立目录) -- Configer `state.judger` — 流水线进度和产出路径 - -如果没有 `judger.pkl`,说明没有走正确路径。 - -## When to Use - -当 Codex 或用户要执行以下操作时使用本 skill: - -- 评测模型生成的代码 / SQL / 文本 -- 计算 pass@k、accuracy 等指标 -- 从 Configer(DB)配置启动评测流水线 -- 断点续跑中断的评测任务 -- 查看评测进度事件及 pass@k/stats - -不要用它处理: - -- 训练、数据爬取、数据构造(走对应的 Agent/Skill) -- 全局 `system` 配置修改(走 Configer) - -## Configuration - -Codex 启动评测前,必须通过 **Configer skill** 检查并预填配置。**不要手动拼字段**——用 `configer_get_task` 查缺,用 `configer_update_task` 补齐。 - -### 预填写流程 +或通过 CLI: -``` -1. configer_get_task(schema="states", section="judger", task_id="") - → 查看字段 schema 和当前 value,找出 value 为 null 的必填字段 -2. 将缺失字段列表和待写入的值告知用户,**必须征得用户确认后才能写入** -3. configer_update_task("judger", {用户确认的字段}, task_id="") - → Configer 会校验字段名是否合法,不存在的字段直接报错 +```bash +DB_PATH=api/db/db.sqlite3 TASK_ID= loopai-judger ``` -**⚠️ 修改 state 前必须询问用户。** 不要自动覆盖已有配置字段,不要猜测 model_path、problem_path 等路径值。 - -### 运行环境 - -| 条件 | 说明 | -|---|---| -| `DB_PATH` | Configer SQLite 数据库(环境变量注入) | -| `TASK_ID` | 任务唯一标识(从环境变量获取) | -| GPU | 至少一张 CUDA GPU | -| Port 8911 | code/text2sql 的 vLLM HTTP 端口 | +## Configuration -### 必填字段 +配置通过 **Configer skill** 写入 `TaskModel.state`,分两部分: -以下字段必须在 `state["judger"]` 中有非空值: +### 全局字段(state["judger"] 顶层,所有 bench 共享) -| 字段 | 适用 task_type | 示例 | +| 字段 | 默认值 | 说明 | |---|---|---| -| `eval_task_type` | 全部 | `"code"` / `"text2sql"` / `"general_text"` | -| `eval_model_path` | 全部 | `"/data/models/Qwen2.5-7B-Instruct/"` | -| `eval_problem_path` | 全部 | `"/data/.../dev_bird_for_oj_sampled.jsonl"` | -| `eval_text2sql_dir` | text2sql | `"/data/.../dev_databases/"` | -| `bench_dataflow_eval_type` | general_text | `"key2_qa"` | - -### 可选字段(有默认值,通常不需要改) - -| 字段 | 默认值 | 何时需要修改 | -|---|---|---| -| `eval_temperature` | `0` | 调整生成随机性 | -| `eval_top_p` | `0.95` | 调整采样策略 | -| `eval_batch_size` | `10` | GPU 显存不足时调小 | -| `eval_case_num` | `10` | 提高 pass@k 精度 | -| `eval_vllm_tensor_parallel_size` | `1` | 多 GPU 时调整 | -| `eval_vllm_gpu_memory_utilization` | `0.9` | GPU 显存不足时调小 | +| `eval_model_path` | 无 | 模型路径(必填) | +| `eval_temperature` | `0` | 采样温度 | +| `eval_top_p` | `0.95` | Top-P 采样 | +| `eval_batch_size` | `10` | 批处理大小,bench 可覆盖 | +| `eval_case_num` | `10` | 每问题样本数,bench 可覆盖 | +| `eval_vllm_tensor_parallel_size` | `1` | vLLM 张量并行数 | +| `eval_vllm_gpu_memory_utilization` | `0.9` | vLLM GPU 显存利用率 | | `cuda_visible_devices` | `"0"` | 指定 GPU | -| `output_dir` | `"./outputs"` | 自定义输出路径 | -| `bench_name` | `"general_text_eval"` | general_text 基准名 | -| `key_mapping` | `{}` | general_text 字段映射(可自动推断) | - -### general_text 评测类型 - -| `bench_dataflow_eval_type` | 说明 | -|---|---| -| `key1_text_score` | 文本评分 | -| `key2_qa` | 问答评测 | -| `key2_q_ma` | 多答案评测 | -| `key3_q_choices_a` | 选择题评测 | -| `key3_q_choices_as` | 多选评测 | -| `key3_q_a_rejected` | 对比评测 | +| `output_dir` | `"./outputs"` | 输出根目录 | -### 配置优先级 - -``` -环境变量 > state["judger"](DB)> schema 默认值 -``` - -### validate 自动检查 - -- 必填字段是否缺失 → `CONFIG_ERROR`,detail 列出 `missing_fields` -- `eval_problem_path` 文件是否存在 → `NOT_FOUND` -- JSONL 字段结构是否匹配 task_type → `INVALID_INPUT` - -``` -loopai/skills/Judger/ ← Skill 层(独立模式,无 LangGraph) -├── __init__.py # run() / load_events() -├── runner.py # 流水线主逻辑 + _load_task_state / _save_task_progress -├── runtime_config.py # 配置解析(env > state["judger"] > schema defaults) -└── utils/ - ├── eval_general_text.py # general_text 评测(One-Eval DataFlowEvalTool) - ├── generate.py # code/text2sql 样本生成 - ├── evaluate.py # code/text2sql 评测(含 pass@k 计算) - └── format.py # 数据格式转换 -``` - -## Python API - -`loopai.skills.Judger.run()` 是唯一入口,全部参数如下: - -```python -from loopai.skills.Judger import run - -result = run( - state=None, # dict with state["judger"] fields,None 时从 Configer 加载 - resume=False, # True = 从 Configer 恢复上次进度 - from_step=None, # 强制起始步骤名 -) -``` - -## Pipeline Steps - -### code / text2sql pipeline - -``` -validate → kill_vllm → start_vllm → format_data → generate → evaluate → kill_vllm_cleanup → finish -``` +### Bench 配置(state["judger"]) -### general_text pipeline - -``` -validate → eval_general_text → finish -``` - - -| Step | 功能 | 完成事件 data | -| ------------------- | ------------------------------- | ----------------------------------------------------- | -| `validate` | 校验必填字段、文件存在性、JSONL 字段结构 | `task_type`, `problem_path` | -| `kill_vllm` | 关闭端口 8911 上的 vLLM 进程 | — | -| `start_vllm` | 启动本地 vLLM 服务 | `base_url` | -| `format_data` | 数据格式转换(human-eval / mbpp),可选 | `target` | -| `generate` | vLLM 批量生成 code/text2sql 样本 | `output_case_path` | -| `evaluate` | 执行代码/执行 SQL,计算 pass@k | `output_result_path`, `**metrics**` | -| `kill_vllm_cleanup` | 评测后关闭 vLLM | — | -| `eval_general_text` | One-Eval DataFlowEvalTool 子进程评测 | `output_result_path`, `output_pred_path`, `**metrics**` | -| `finish` | 流水线完成 | — | - - -## Output Metrics - -Judger 在不同任务类型下产出的评测指标,均写入事件流(`judger.pkl`)和输出文件。 - -### code / text2sql — pass@k - -由 Judger 直接计算(`evaluate.py` → `_calculate_pass_at_k`),不需要 One-Eval。 - -| 指标 | 说明 | 计算方式 | -|---|---|---| -| `pass@1` | 1 次采样通过率 | `estimate_pass_at_k(n, c, 1)` | -| `pass@10` | 10 次采样通过率(需 `eval_case_num ≥ 10`) | `estimate_pass_at_k(n, c, 10)` | -| `pass@100` | 100 次采样通过率(需 `eval_case_num ≥ 100`) | `estimate_pass_at_k(n, c, 100)` | - -- 输出位置:事件流 `data.metrics` + `outputs//judger//log.txt` -- k 值列表:`[1, 10, 100]`,仅当 `total_samples ≥ k` 时对应 k 才会出现在结果中 - -**事件示例:** +所有评测集通过 `benchlist` 和 `extra_benchlist` 列表配置。**格式必须是 JSON 数组**(`[{...},{...}]`),**不是** JSONL(每行一个对象): ```json -{ - "current": "judger.evaluate", - "progress": 1.0, - "message": "评测完成", - "data": { - "output_result_path": "outputs/.../result.jsonl", - "metrics": { - "pass@1": 0.3125 - } - } -} +[{"name":"gsm8k","task_type":"general_text","problem_path":"/data/gsm8k/test.jsonl","eval_type":"key2_qa"},{"name":"human_eval","task_type":"code","problem_path":"/data/humaneval.jsonl","case_num":10}] ``` -### general_text — One-Eval stats - -由 One-Eval `DataFlowEvalTool` 计算并返回 `stats` 字典。具体包含哪些指标取决于 `bench_dataflow_eval_type`。 - -**所有 eval_type 通用的 stat 键:** - -| 键 | 类型 | 说明 | -|---|---|---| -| `accuracy` | `float` | 综合准确率(0~1) | -| `score` | `float` | 综合得分(通常 = accuracy) | -| `total_samples` | `int` | 总样本数 | -| `valid_samples` | `int` | 有效样本数 | - -**按 eval_type 的专属指标:** - -| `bench_dataflow_eval_type` | 典型产出指标 | -|---|---| -| `key1_text_score` | `bleu`, `rouge`, `chrf`, `ter`, `token_f1`, `exact_match`, `containment_match` | -| `key2_qa` | `exact_match`, `containment_match`, `numerical_match`, `token_f1` | -| `key2_q_ma` | `exact_match`, `token_f1` | -| `key3_q_choices_a` | `choice_accuracy`, `exact_match` | -| `key3_q_choices_as` | `exact_match`, `token_f1` | -| `key3_q_a_rejected` | 对比评测指标(pairwise comparison) | - -**One-Eval 支持的完整指标集:** - -| 指标名 | 类别 | 说明 | -|---|---|---| -| `pass_at_k` | code | 代码 pass@k(One-Eval 版本) | -| `code_similarity` | code | 代码相似度 | -| `soft_code_execution` | code | 软代码执行评测 | -| `exact_match` | general | 精确匹配 | -| `containment_match` | general | 包含匹配 | -| `strict_match` | general | 严格匹配 | -| `numerical_match` | general | 数值匹配 | -| `choice_accuracy` | general | 选择题准确率 | -| `bleu` | text_gen | BLEU 机器翻译评测 | -| `rouge` | text_gen | ROUGE 摘要评测 | -| `chrf` | text_gen | 字符级 n-gram F-score | -| `ter` | text_gen | 翻译错误率 | -| `token_f1` | text_gen | Token 级 F1 | -| `math_verify` | math | 数学表达式验证 | -| `symbolic_match` | symbolic | 符号匹配 | -| `spearman` | classification | Spearman 排名相关系数 | -| `pearson` | classification | Pearson 相关系数 | -| `mcc` | classification | Matthews 相关系数 | -| `auc_roc` | classification | AUC-ROC | -| `gini_index` | classification | Gini 系数 | - -> **注意**:上表为 One-Eval 的完整能力。实际 stats 中出现的指标由 One-Eval 根据 task_type 和数据特征自动选择。不是所有指标都会同时出现。 - -**事件示例:** - ```json { - "current": "judger.eval_general_text", - "progress": 1.0, - "message": "通用文本评测完成", - "data": { - "output_result_path": "outputs/.../text_eval_summary_20240625_120000.json", - "output_pred_path": "outputs/.../text_eval_scored_20240625_120000.json", - "metrics": { - "accuracy": 0.94, - "score": 0.94, - "total_samples": 100, - "valid_samples": 95, - "token_f1": 0.89 + "benchlist": [ + { + "name": "gsm8k", + "task_type": "general_text", + "problem_path": "/data/gsm8k/test.jsonl", + "eval_type": "key2_qa", + "key_mapping": {} + }, + { + "name": "human_eval", + "task_type": "code", + "problem_path": "/data/humaneval.jsonl", + "case_num": 10, + "batch_size": 10, + "format_type": "" + }, + { + "name": "bird_dev", + "task_type": "text2sql", + "problem_path": "/data/bird/dev.jsonl", + "text2sql_dir": "/data/bird/dev_databases", + "case_num": 10, + "batch_size": 10 } - } + ], + "extra_benchlist": [] } ``` -### 指标产出总结 - -| 指标类型 | code/text2sql | general_text | 输出位置 | -|---|---|---|---| -| `pass_at_k` | ✅ | — | stdout (metrics) + 事件流 + log.txt | -| `accuracy` | — | ✅ | stdout (metrics) + 事件流 + summary JSON | -| `score` | — | ✅ | stdout (metrics) + 事件流 + summary JSON | -| `token_f1` | — | ✅ | stdout (metrics) + 事件流 + summary JSON | -| `exact_match` | — | ✅ | stdout (metrics) + 事件流 + summary JSON | -| `bleu` / `rouge` / `chrf` | — | ✅ (text_score) | stdout (metrics) + 事件流 + summary JSON | -| `spearman` (ranking) | — | ✅ (classification) | stdout (metrics) + 事件流 + summary JSON | -| `reward_score` | — | ❌ (One-Eval 不直接产出) | — | - -### 如何读取指标 - -```python -from loopai.skills.Judger import load_events - -events = load_events(task_id="my_task") -for e in events: - data = e.get("data") or {} - - # 统一指标:code/text2sql → pass@k,general_text → accuracy/score/f1 等 - if "metrics" in data: - for k, v in data["metrics"].items(): - print(f"{k}: {v:.4f}") -``` - -## Environment Variables - - -| 变量 | 对应字段 | 默认值 | -| --------------------------------- | ---------------------------------- | ------------------- | -| `DB_PATH` | Configer 数据库路径 | 必填(Configer 模式) | -| `TASK_ID` | `task_id` | 必填,无默认 | -| `OUTPUT_DIR` | `output_dir` | `./outputs` | -| `JUDGER_MODEL_PATH` | `eval_model_path` | 必填(无 DB/YAML 时) | -| `JUDGER_TASK_TYPE` | `eval_task_type` | `code` | -| `JUDGER_TEMPERATURE` | `eval_temperature` | `0` | -| `JUDGER_TOP_P` | `eval_top_p` | `0.95` | -| `JUDGER_PROBLEM_PATH` | `eval_problem_path` | 必填(无 DB/YAML 时) | -| `JUDGER_BATCH_SIZE` | `eval_batch_size` | `10` | -| `JUDGER_CASE_NUM` | `eval_case_num` | `10` | -| `JUDGER_FORMAT_TYPE` | `eval_format_type` | 可选 | -| `JUDGER_TEXT2SQL_DIR` | `eval_text2sql_dir` | text2sql 必填 | -| `JUDGER_TENSOR_PARALLEL_SIZE` | `eval_vllm_tensor_parallel_size` | `1` | -| `JUDGER_GPU_MEMORY_UTILIZATION` | `eval_vllm_gpu_memory_utilization` | `0.9` | -| `CUDA_VISIBLE_DEVICES` | `cuda_visible_devices` | `0` | -| `JUDGER_BENCH_NAME` | `bench_name` | `general_text_eval` | -| `JUDGER_BENCH_DATAFLOW_EVAL_TYPE` | `bench_dataflow_eval_type` | 空(general_text 必填) | - - -## Output & Artifacts - -``` -outputs// -├── judger/ -│ └── / # 每次运行独立目录 -│ ├── _sample.jsonl # 生成的样本 -│ ├── _result.jsonl # 评测结果 -│ ├── log.txt # 评测日志(含 pass@k) -│ ├── text_eval_summary_*.json # general_text 摘要 -│ ├── general_text_dataset_cache_*.jsonl # 缓存 -│ └── gsm8k_*_steps/ # One-Eval 中间产物 -└── judger.pkl # 事件 pickle(load_events 读取) -``` - -- **stdout** — 最终结果 JSON payload(`emit_success` / `emit_error`,Codex 消费) -- **judger.pkl** — 所有进度事件,含步骤完成时的 `metrics`(pass@k 或 stats) -- **log.txt** — 评测日志(`outputs//judger//log.txt`),含 pass@k 数值 - -## State & Resume(Configer) - -### 工作原理 - -Judger 通过 Configer 读写 TaskModel.state: - -- **读取**:`_load_task_state(task_id)` 调用 `get_configer_task_state_config` 从 DB 读取 `state.judger` -- **写入**:每步执行前后调用 `_save_task_progress` → `update_configer_task_state_config` 写入进度 - -state 中的关键字段: - - -| 字段 | 用途 | -| --------------------------------- | ------------------------- | -| `state.judger._last_completed` | 最后完成的步骤名(如 `evaluate`) | -| `state.judger._current` | 当前步骤(如 `judger.generate`) | -| `state.judger.output_result_path` | 评测结果路径 | -| `state.judger.output_case_path` | 样本路径 | - - -### 断点续跑 - -```python -from loopai.skills.Judger import run - -# 从上次中断处继续 -run(resume=True) - -# 从指定步骤强制执行 -run(from_step="evaluate") -``` - -`resume` 时 state 从 Configer 加载;`from_step` 指定起始步骤跳过之前所有步骤。 - -`_is_finished` 检查:如果 `last_completed == "finish"`,流水线跳过所有步骤直接返回。 +**bench entry 字段:** -## Event System +| 字段 | code | text2sql | general_text | 说明 | +|---|---|---|---|---| +| `name` | ✅ 必填 | ✅ 必填 | ✅ 必填 | bench 标识 | +| `task_type` | ✅ 必填 | ✅ 必填 | ✅ 必填 | `code` / `text2sql` / `general_text` | +| `problem_path` | ✅ 必填 | ✅ 必填 | ✅ 必填 | 问题文件路径 | +| `case_num` | 可选 10 | 可选 10 | — | 每问题样本数,bench 设了覆盖全局 | +| `batch_size` | 可选 10 | 可选 10 | — | 批处理大小,bench 设了覆盖全局 | +| `format_type` | 可选 | — | — | `human-eval` / `mbpp`,不设走默认 | +| `text2sql_dir` | — | ✅ 必填 | — | SQLite 数据库目录 | +| `eval_type` | — | — | ✅ 必填 | `key2_qa` / `key1_text_score` 等 | +| `key_mapping` | — | — | 可选 | 字段映射,可自动推断 | -### 事件写入 +**主/附加区别:** -流水线运行时自动持久化到 `//judger.pkl`: - -```python -from loopai.common.event_tool import get_event_writer, StreamEvent - -writer = get_event_writer(name="judger", context_id="task_001", log_file_path="./outputs") -writer(StreamEvent(current="judger.generate", progress=0.5, message="样本生成中")) -``` - -### 事件格式 - -每个事件包含: - -- `current` — 当前步骤(格式 `judger.`) -- `progress` — 步骤内进度 0.0 ~ 1.0 -- `message` — 人类可读描述 -- `data` — 结构化数据(步骤完成时包含关键结果) -- `status` — 仅终态事件:`"completed"`(成功)或 `"failed"`(失败) - -### 终态事件 - -流水线结束时自动写入: - -```json -// 成功 — writer.set_completed() -{"current": "judger", "status": "completed", "message": "Sub-agent completed."} - -// 失败 — emit_error(stream_writer=writer) → writer.set_failed() -{"current": "judger", "status": "failed", "message": "Sub-agent failed.", "error": {...}} -``` - -### 步骤完成事件 data +| | 主任务 | 附加任务 | +|---|---|---| +| 执行顺序 | 先 | 后 | +| 失败策略 | 记录失败 + `_save_task_progress` + 退出 | 记录失败,继续 | -evaluate 步骤完成时(code/text2sql): +### 预填写流程 -```json -{ - "output_result_path": "outputs/.../result.jsonl", - "metrics": { - "pass@1": 0.85, - "pass@10": 0.95, - "pass@100": 1.0 - } -} ``` - -eval_general_text 步骤完成时: - -```json -{ - "output_result_path": "outputs/.../summary.json", - "output_pred_path": "outputs/.../step2.jsonl", - "metrics": { - "accuracy": 0.94, - "score": 0.88, - "total_samples": 100, - "valid_samples": 95 - } -} +1. configer_get_task(schema="states", section="judger", task_id="") +2. 将缺失字段告知用户,征得确认后写入 +3. configer_update_task("judger", {"benchlist": [...], "eval_model_path": "..."}, task_id="") ``` -### 事件读取 +## Pipeline -```python -from loopai.skills.Judger import load_events +每个 bench entry 独立跑一遍完整流水线: -events = load_events(task_id="my_task") -for e in events: - data = e.get("data") or {} - if "metrics" in data: - for k, v in data["metrics"].items(): - print(f"{k}: {v:.4f}") ``` - -## Error Handling - -每个步骤失败点直接调用 `emit_error(exc, stream_writer=writer)`,**三通道同时输出**: - -| 通道 | 机制 | 内容 | -|---|---|---| -| stdout | `print` error JSON | `{"ok": false, "error": {"code": "CONFIG_ERROR", ...}}` | -| judger.pkl | `stream_writer.set_failed()` → `_append_status_event` | status="failed" 事件 | -| DB | `stream_writer._sync_runtime(status="failed")` | taskruntime 表标记失败 | - -成功时 pipeline 末尾调 `writer.set_completed()`,同样写 judger.pkl + DB。 - -| ErrorCode | 触发场景 | -|---|---| -| `CONFIG_ERROR` | 缺少必填字段、模型路径未配置、task_id 缺失 | -| `INVALID_INPUT` | JSONL 字段不匹配、不支持的任务类型、未知步骤名 | -| `NOT_FOUND` | 问题文件不存在 | -| `EXTERNAL_SERVICE_ERROR` | DataFlowEvalTool 子进程失败、vLLM 启动失败 | -| `UNHANDLED_EXCEPTION` | 意外的未分类异常 | - -### 错误响应格式(stdout) - -```json -{ - "ok": false, - "status": "failed", - "message": "Judger configuration is incomplete.", - "data": null, - "error": { - "type": "ValueError", - "code": "CONFIG_ERROR", - "detail": "Missing required fields: ...", - "traceback": "...", - "recoverable": true, - "time": "2026-06-25T12:00:00Z" - } -} +对每个 bench: + _apply_bench_to_state → 注入 bench 字段到 state["judger"] + → 按 task_type 选流水线: + code/text2sql: validate → kill_vllm → start_vllm → format_data → generate → evaluate → kill_vllm_cleanup → finish + general_text: validate → eval_general_text → finish + → 收集结果到 bench_result / extra_bench_result ``` -### 成功响应格式 +## Output -**code/text2sql:** +### stdout(emit_success) ```json { - "ok": true, - "status": "completed", - "message": "Judger pipeline completed.", - "data": { - "task_type": "text2sql", - "output_result_path": "/path/to/result.jsonl", - "output_case_path": "/path/to/sample.jsonl", - "output_problem_path": "/path/to/problem.jsonl", - "output_pred_path": "", - "bench": "", - "metrics": { - "pass@1": 0.3125 - } - }, - "error": null + "ok": true, + "data": { + "bench_result": [ + {"bench_name": "gsm8k", "task_type": "general_text", + "output_result_path": "...", "metrics": {"accuracy": 0.94}} + ], + "extra_bench_result": [ + {"bench_name": "human_eval", "task_type": "code", + "output_result_path": "...", "metrics": {"pass@1": 0.85}} + ], + "metrics": {"gsm8k": {"accuracy": 0.94}, "human_eval": {"pass@1": 0.85}} + } } ``` -**general_text:** +### 目录结构 -```json -{ - "ok": true, - "status": "completed", - "message": "Judger pipeline completed.", - "data": { - "task_type": "general_text", - "output_result_path": "/path/to/summary.json", - "output_case_path": "", - "output_problem_path": "/path/to/cache.jsonl", - "output_pred_path": "/path/to/scored.jsonl", - "bench": { - "bench_name": "gsm8k", - "eval_status": "success", - "meta": {"eval_result": {"accuracy": 0.94}} - }, - "metrics": { - "accuracy": 0.94, - "score": 0.94, - "total_samples": 100, - "valid_samples": 95 - } - }, - "error": null -} ``` - -## vLLM Management - -**仅支持本地 vLLM 启动。** 远程 API(`eval_base_url`)在独立模式下不支持。 - -流水线自动管理 vLLM 生命周期: - -1. 关闭端口 8911 上已有的 vLLM 进程 -2. 使用配置的 model_path、tensor_parallel_size、gpu_memory_utilization 启动 vLLM -3. code/text2sql 评测完成后自动关闭;general_text 由 One-Eval 自行管理 - -## Codex Integration - -Codex 通过子进程调用 Python 函数,读取 stdout JSON: - -```bash -timeout 600 python3 -u <<'PY' -import json, os, sys -os.environ["DB_PATH"] = "api/db/db.sqlite3" -os.environ["TASK_ID"] = "codex_task_001" - -from loopai.skills.Judger import run - -try: - result = run( - state={ - "judger": { - "eval_model_path": "/data/models/Qwen2.5-7B-Instruct", - "eval_task_type": "general_text", - "eval_problem_path": "/data/test.jsonl", - "bench_dataflow_eval_type": "key2_qa", - "eval_batch_size": 4, - "cuda_visible_devices": "5", - }, - "output_dir": "./outputs", - }, - ) - # 成功:emit_success 输出到 stdout - sys.exit(0) -except Exception as e: - # 错误:emit_error 已输出到 stdout - sys.exit(1) -PY +outputs// +├── judger/ +│ └── / +│ ├── gsm8k/ ← bench_name 子目录 +│ │ ├── text_eval_summary_*.json +│ │ └── gsm8k_*_steps/ +│ ├── human_eval/ +│ │ ├── human_eval_sample.jsonl +│ │ ├── human_eval_result.jsonl +│ │ └── log.txt +│ └── bird_dev/ +└── judger.pkl ``` -Codex 读到的 stdout 行: - -- 成功:`{"ok": true, "data": {"output_result_path": "...", "bench": {...}}}` -- 失败:`{"ok": false, "error": {"code": "CONFIG_ERROR", "detail": "..."}}` - -Codex 也可以通过 `load_events` 读取 judger.pkl 获取 pass@k / stats 等评测指标。 +### Configer 持久化 -## Config Via Configer +`_save_task_progress` 写入 `state.judger.bench_result` 和 `state.judger.extra_bench_result`,Analyzer 从中读取。 -Judger 配置字段可通过 Configer skill 读写(按 task_id 隔离): - -```python -from loopai.skills.Configer import ( - get_configer_state_schema, - get_configer_task_state_config, - update_configer_task_state_config, -) +## Error Handling -# 查看 schema(字段含义、允许值、默认值) -schema = get_configer_state_schema(section_name="judger") +每个步骤 `emit_error(exc, stream_writer=writer)`: +- stdout 输出 `{"ok": false, ...}` +- judger.pkl 写入 `status=failed` +- taskruntime 表标记失败 -# 读取某个 task 的当前配置 -config = get_configer_task_state_config( - section_name="judger", - task_id="a4341a82-4ed4-46da-8776-d9cf45a4f50c", -) +所有 error `recoverable=true`,Codex 可引导用户修复后重试。 -# 修改某个 task 的配置(运行前预设参数) -update_configer_task_state_config( - "judger", - {"eval_temperature": 0.2, "eval_case_num": 20}, - task_id="a4341a82-4ed4-46da-8776-d9cf45a4f50c", -) -``` +## Environment Variables -**流水线进度也通过 Configer 持久化**:`_last_completed` 和 `_current` 字段在每步前后自动写入 `state.judger`,resume 时从中恢复。 +| 变量 | 来源 | 默认值 | +|---|---|---| +| `DB_PATH` | 环境变量 | 必填 | +| `TASK_ID` | 环境变量 | 必填 | +| `OUTPUT_DIR` | 环境变量 | `./outputs` | +| `CUDA_VISIBLE_DEVICES` | 环境变量 | `"0"` | diff --git a/tests/test_judger.py b/tests/test_judger.py new file mode 100644 index 0000000..6f0b7f8 --- /dev/null +++ b/tests/test_judger.py @@ -0,0 +1,107 @@ +#!/usr/bin/env python +# -*- coding: utf-8 -*- +"""Judger 集成测试 — 传入 bench 配置,跑完整 pipeline。 + +用法: + conda activate loopai_zx + DB_PATH=api/db/db.sqlite3 TASK_ID= python tests/test_judger.py \ + --bench-config bench_config.json + +bench_config.json 格式: +{ + "benchlist": [ + {"name": "gsm8k", "task_type": "general_text", + "problem_path": "/data/gsm8k/test.jsonl", "eval_type": "key2_qa"} + ], + "extra_benchlist": [ + {"name": "mmlu", "task_type": "general_text", + "problem_path": "/data/mmlu/test.jsonl", "eval_type": "key3_q_choices_a"} + ] +} +""" + +from __future__ import annotations + +import argparse +import json +import os +import sys +from pathlib import Path + +_PROJECT_ROOT = Path(__file__).resolve().parent.parent +if str(_PROJECT_ROOT) not in sys.path: + sys.path.insert(0, str(_PROJECT_ROOT)) + + +def load_bench_config(path: str) -> dict: + with open(path, "r", encoding="utf-8") as f: + return json.load(f) + + +def main(): + parser = argparse.ArgumentParser(description="Judger 集成测试") + default_config = str(Path(__file__).parent / "bench_config.json") + parser.add_argument("--bench-config", type=str, default=default_config, + help=f"Bench 配置 JSON 文件路径(默认: {default_config})") + parser.add_argument("--dry-run", action="store_true", + help="只校验不走流水线") + args = parser.parse_args() + + db_path = os.getenv("DB_PATH", "api/db/db.sqlite3") + task_id = os.getenv("TASK_ID") + if not task_id: + print("❌ 请设置 TASK_ID 环境变量", file=sys.stderr) + sys.exit(1) + + # 加载 bench 配置 + benches = load_bench_config(args.bench_config) + primary = benches.get("benchlist", []) + secondary = benches.get("extra_benchlist", []) + + if not primary and not secondary: + print("❌ bench_config.json 中 benchlist 或 extra_benchlist 至少需要一个", file=sys.stderr) + sys.exit(1) + + print(f"benchlist: {[b.get('name') for b in primary]}") + print(f"extra_benchlist: {[b.get('name') for b in secondary]}") + + if args.dry_run: + # 只校验 bench entry 结构 + from loopai.skills.Judger.runner import _validate_bench + all_ok = True + for b in primary + secondary: + try: + _validate_bench(b, writer=None) + print(f" ✅ {b.get('name')}") + except SystemExit: + print(f" ❌ {b.get('name')} 校验失败") + all_ok = False + if all_ok: + print("\n✅ 所有 bench 校验通过") + else: + sys.exit(1) + return + + # 写 Configer:如果 task 还没配置,写入模型路径和 benches + # (模型路径从 state 或 env 读取) + from loopai.skills.Configer import update_configer_task_state_config + updates = { + "benchlist": primary, + "extra_benchlist": secondary, + } + model_path = os.getenv("JUDGER_MODEL_PATH") + if model_path: + updates["eval_model_path"] = model_path + update_configer_task_state_config("judger", updates, task_id=task_id) + print("✅ bench 配置已写入 Configer") + + # 跑 pipeline + os.environ["DB_PATH"] = db_path + os.environ["TASK_ID"] = task_id + from loopai.skills.Judger import run + print("🚀 启动 Judger pipeline...\n") + run() + + +if __name__ == "__main__": + main()