#!/usr/bin/env python3 """Local EvalScope launch panel — wraps bash/run.py via FastAPI.""" from __future__ import annotations import asyncio import json import os import signal import sys import uuid from datetime import datetime, timezone from pathlib import Path from typing import Any, Dict, List, Optional from fastapi import FastAPI, HTTPException from fastapi.responses import FileResponse, StreamingResponse from fastapi.staticfiles import StaticFiles from pydantic import BaseModel, Field WEBUI_DIR = Path(__file__).parent.resolve() PROJECT_ROOT = WEBUI_DIR.parent BASH_DIR = PROJECT_ROOT / 'bash' RUN_SCRIPT = BASH_DIR / 'run.py' STATIC_DIR = WEBUI_DIR / 'static' DATA_DIR = WEBUI_DIR / 'data' JOBS_DIR = DATA_DIR / 'jobs' LOGS_DIR = DATA_DIR / 'logs' sys.path.insert(0, str(BASH_DIR)) import run as run_module # noqa: E402 import results_scan # noqa: E402 JOBS_DIR.mkdir(parents=True, exist_ok=True) LOGS_DIR.mkdir(parents=True, exist_ok=True) DEFAULT_OUTPUT_DIR = Path(run_module.DEFAULT_OUTPUT_DIR) app = FastAPI(title='EvalStone Launch Panel', version='1.0.0') # --------------------------------------------------------------------------- # Job store # --------------------------------------------------------------------------- class JobRecord: def __init__(self, job_id: str, payload: dict, command: List[str]): self.id = job_id self.payload = payload self.command = command self.status = 'queued' # queued | running | completed | failed | stopped self.created_at = datetime.now(timezone.utc).isoformat() self.started_at: Optional[str] = None self.finished_at: Optional[str] = None self.return_code: Optional[int] = None self.pid: Optional[int] = None self.log_path = LOGS_DIR / f'{job_id}.log' self.error: Optional[str] = None self._proc: Optional[asyncio.subprocess.Process] = None self._log_fp = None def to_dict(self) -> dict: return { 'id': self.id, 'status': self.status, 'created_at': self.created_at, 'started_at': self.started_at, 'finished_at': self.finished_at, 'return_code': self.return_code, 'pid': self.pid, 'command': self.command, 'payload': self.payload, 'log_path': str(self.log_path), 'error': self.error, } def save(self) -> None: path = JOBS_DIR / f'{self.id}.json' path.write_text(json.dumps(self.to_dict(), ensure_ascii=False, indent=2), encoding='utf-8') JOBS: Dict[str, JobRecord] = {} _ACTIVE_JOB_ID: Optional[str] = None _LOCK = asyncio.Lock() def _load_existing_jobs() -> None: for path in sorted(JOBS_DIR.glob('*.json'), key=lambda p: p.stat().st_mtime, reverse=True): try: data = json.loads(path.read_text(encoding='utf-8')) job = JobRecord(data['id'], data.get('payload', {}), data.get('command', [])) job.status = data.get('status', 'unknown') job.created_at = data.get('created_at', job.created_at) job.started_at = data.get('started_at') job.finished_at = data.get('finished_at') job.return_code = data.get('return_code') job.pid = data.get('pid') job.error = data.get('error') if job.status == 'running': # Process cannot be resumed after server restart job.status = 'failed' job.error = 'Server restarted while job was running' job.finished_at = datetime.now(timezone.utc).isoformat() job.save() JOBS[job.id] = job except Exception: continue _load_existing_jobs() # --------------------------------------------------------------------------- # Request models # --------------------------------------------------------------------------- class LaunchRequest(BaseModel): model: str = Field(..., min_length=1) api_url: str = Field(..., min_length=1) api_key: str = 'EMPTY' thinking: bool = False selection_mode: str = 'suite' # suite | datasets suite: str = 'official' datasets: List[str] = Field(default_factory=list) exclude: List[str] = Field(default_factory=list) folder_name: Optional[str] = None limit: Optional[str] = None seed: int = 42 batch_size: int = 4 thinking_max_tokens_scale: float = 1.0 max_tokens_add: int = 0 dataset_dir: Optional[str] = None output_dir: Optional[str] = None config: Optional[str] = None tokenizer_path: Optional[str] = None judge_model: Optional[str] = None judge_api_url: Optional[str] = None judge_api_key: Optional[str] = None judge_max_tokens: Optional[int] = None write_summary: bool = True # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- # Capability-domain categories for the custom benchmark picker. # Order here is the display order in the UI. BENCHMARK_CATEGORIES = [ { 'id': 'math', 'name': '数学推理', 'items': [ 'aime24', 'aime25', 'aime26', 'hmmt26', 'imo_answerbench', 'competition_math', 'gsm8k', ], }, { 'id': 'code', 'name': '代码', 'items': ['humaneval', 'live_code_bench', 'bigcodebench'], }, { 'id': 'science', 'name': '科学 / 高难推理', 'items': ['gpqa_diamond', 'super_gpqa', 'hle'], }, { 'id': 'knowledge', 'name': '知识与通用能力', 'items': [ 'mmlu', 'mmlu_pro', 'cmmlu', 'bbh', 'arc', 'drop', 'hellaswag', 'winogrande', 'simple_qa', 'trivia_qa', ], }, { 'id': 'long_context', 'name': '长文本', 'items': ['longbench_v2', 'openai_mrcr'], }, { 'id': 'tool_agent', 'name': '工具调用 / 智能体', 'items': ['bfcl_v3', 'general_fc', 'tau2_bench'], }, ] def _meta() -> dict: all_benchmarks = sorted( set(run_module.ALL_MULTI_RUN) | set(run_module.ALL_SINGLE_RUN) | set(run_module.ALL_AGENT) ) categorized = {b for cat in BENCHMARK_CATEGORIES for b in cat['items']} other = sorted(set(all_benchmarks) - categorized) categories = [dict(cat) for cat in BENCHMARK_CATEGORIES] if other: categories.append({'id': 'other', 'name': '其他', 'items': other}) suites = {} for name, cfg in run_module.SUITES.items(): suites[name] = { 'multi': list(cfg['multi']), 'single': list(cfg['single']), 'agent': list(cfg['agent']), 'all': list(cfg['multi']) + list(cfg['single']) + list(cfg['agent']), } return { 'suites': suites, 'benchmarks': all_benchmarks, 'categories': categories, 'multi_run': run_module.MULTI_RUN_CONFIG, 'defaults': { 'model': run_module.DEFAULT_MODEL, 'api_url': run_module.DEFAULT_API_URL, 'api_key': 'EMPTY', 'dataset_dir': run_module.DEFAULT_DATASET_DIR, 'output_dir': run_module.DEFAULT_OUTPUT_DIR, 'config': run_module.DEFAULT_CONFIG, 'tokenizer_path': run_module.DEFAULT_TOKENIZER_PATH, 'seed': run_module.DEFAULT_SEED, 'batch_size': run_module.DEFAULT_BATCH_SIZE, 'thinking': run_module.DEFAULT_ENABLE_THINKING, 'judge_model': run_module.DEFAULT_JUDGE_MODEL, 'judge_api_url': run_module.DEFAULT_JUDGE_API_URL, 'judge_max_tokens': run_module.DEFAULT_JUDGE_MAX_TOKENS, 'suite': 'official', }, 'project_root': str(PROJECT_ROOT), 'run_script': str(RUN_SCRIPT), } def build_command(req: LaunchRequest) -> List[str]: cmd = [ sys.executable, str(RUN_SCRIPT), '--model', req.model, '--api-url', req.api_url, '--seed', str(req.seed), '--batch-size', str(req.batch_size), ] # Keep API key in env only; avoid requiring a custom --api-key CLI flag in run.py. if req.thinking: cmd.append('--thinking') else: cmd.append('--no-thinking') if req.selection_mode == 'datasets': if not req.datasets: raise HTTPException(status_code=400, detail='请至少选择一个 benchmark') cmd.extend(['--datasets', ','.join(req.datasets)]) else: if req.suite not in run_module.SUITES: raise HTTPException(status_code=400, detail=f'未知 suite: {req.suite}') cmd.extend(['--suite', req.suite]) if req.exclude: cmd.extend(['--exclude', ','.join(req.exclude)]) if req.folder_name: cmd.extend(['--folder-name', req.folder_name]) if req.limit is not None and str(req.limit).strip() != '': cmd.extend(['--limit', str(req.limit)]) if req.thinking_max_tokens_scale != 1.0: cmd.extend(['--thinking-max-tokens-scale', str(req.thinking_max_tokens_scale)]) if req.max_tokens_add: cmd.extend(['--max-tokens-add', str(req.max_tokens_add)]) if req.dataset_dir: cmd.extend(['--dataset-dir', req.dataset_dir]) if req.output_dir: cmd.extend(['--output-dir', req.output_dir]) if req.config: cmd.extend(['--config', req.config]) if req.tokenizer_path: cmd.extend(['--tokenizer-path', req.tokenizer_path]) if req.judge_model: cmd.extend(['--judge-model', req.judge_model]) if req.judge_api_url: cmd.extend(['--judge-api-url', req.judge_api_url]) if req.judge_api_key: cmd.extend(['--judge-api-key', req.judge_api_key]) if req.judge_max_tokens is not None: cmd.extend(['--judge-max-tokens', str(req.judge_max_tokens)]) if not req.write_summary: cmd.append('--no-summary') return cmd async def _pump_stdout(job: JobRecord) -> None: assert job._proc is not None and job._log_fp is not None assert job._proc.stdout is not None try: while True: line = await job._proc.stdout.readline() if not line: break text = line.decode('utf-8', errors='replace') if job._log_fp and not job._log_fp.closed: job._log_fp.write(text) job._log_fp.flush() return_code = await job._proc.wait() job.return_code = return_code job.finished_at = datetime.now(timezone.utc).isoformat() if job.status in ('stopping', 'stopped'): job.status = 'stopped' job.error = job.error or 'Stopped by user' elif job.status == 'running': if return_code == 0: job.status = 'completed' else: job.status = 'failed' job.error = f'Process exited with code {return_code}' finally: job.pid = None job._proc = None if job._log_fp and not job._log_fp.closed: try: job._log_fp.close() except Exception: pass job._log_fp = None job.save() global _ACTIVE_JOB_ID if _ACTIVE_JOB_ID == job.id: _ACTIVE_JOB_ID = None async def start_job(req: LaunchRequest) -> JobRecord: global _ACTIVE_JOB_ID async with _LOCK: if _ACTIVE_JOB_ID and _ACTIVE_JOB_ID in JOBS and JOBS[_ACTIVE_JOB_ID].status == 'running': raise HTTPException(status_code=409, detail=f'已有任务在运行: {_ACTIVE_JOB_ID}') cmd = build_command(req) job_id = datetime.now().strftime('%Y%m%d_%H%M%S') + '_' + uuid.uuid4().hex[:8] job = JobRecord(job_id, req.model_dump(), cmd) job.log_path.write_text('', encoding='utf-8') env = os.environ.copy() env['PYTHONUNBUFFERED'] = '1' if req.api_key and req.api_key != 'EMPTY': env['OPENAI_API_KEY'] = req.api_key try: proc = await asyncio.create_subprocess_exec( *cmd, cwd=str(PROJECT_ROOT), stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.STDOUT, env=env, start_new_session=True, ) except Exception as e: job.status = 'failed' job.error = str(e) job.finished_at = datetime.now(timezone.utc).isoformat() job.save() JOBS[job.id] = job raise HTTPException(status_code=500, detail=f'启动失败: {e}') from e job._proc = proc job.pid = proc.pid job.status = 'running' job.started_at = datetime.now(timezone.utc).isoformat() job._log_fp = open(job.log_path, 'a', encoding='utf-8') header = ( f'# job {job.id}\n' f'# cwd: {PROJECT_ROOT}\n' f'# cmd: {" ".join(cmd)}\n' f'# started: {job.started_at}\n' f'{"=" * 60}\n' ) job._log_fp.write(header) job._log_fp.flush() job.save() JOBS[job.id] = job _ACTIVE_JOB_ID = job.id asyncio.create_task(_pump_stdout(job)) return job async def stop_job(job_id: str) -> JobRecord: job = JOBS.get(job_id) if not job: raise HTTPException(status_code=404, detail='任务不存在') if job.status != 'running' or job._proc is None: raise HTTPException(status_code=400, detail='任务未在运行') proc = job._proc job.status = 'stopping' job.error = 'Stopped by user' if job._log_fp and not job._log_fp.closed: try: job._log_fp.write('\n# stop requested by user\n') job._log_fp.flush() except Exception: pass job.save() try: os.killpg(proc.pid, signal.SIGTERM) except ProcessLookupError: pass except Exception: proc.terminate() try: await asyncio.wait_for(proc.wait(), timeout=15) except asyncio.TimeoutError: try: os.killpg(proc.pid, signal.SIGKILL) except Exception: try: proc.kill() except Exception: pass # Final status is finalized by _pump_stdout; wait briefly for it. for _ in range(20): if job.status in ('stopped', 'failed', 'completed'): break await asyncio.sleep(0.1) if job.status == 'stopping': job.status = 'stopped' job.finished_at = datetime.now(timezone.utc).isoformat() job.save() return job # --------------------------------------------------------------------------- # API routes # --------------------------------------------------------------------------- @app.get('/api/health') async def health(): return {'ok': True, 'project_root': str(PROJECT_ROOT)} @app.get('/api/meta') async def meta(): return _meta() @app.get('/api/jobs') async def list_jobs(limit: int = 50): items = sorted(JOBS.values(), key=lambda j: j.created_at, reverse=True)[:limit] return {'jobs': [j.to_dict() for j in items], 'active_job_id': _ACTIVE_JOB_ID} @app.get('/api/jobs/{job_id}') async def get_job(job_id: str): job = JOBS.get(job_id) if not job: raise HTTPException(status_code=404, detail='任务不存在') return job.to_dict() @app.post('/api/jobs') async def create_job(req: LaunchRequest): job = await start_job(req) return job.to_dict() @app.post('/api/jobs/{job_id}/stop') async def api_stop_job(job_id: str): job = await stop_job(job_id) return job.to_dict() @app.get('/api/jobs/{job_id}/logs') async def get_logs(job_id: str, offset: int = 0): job = JOBS.get(job_id) if not job: raise HTTPException(status_code=404, detail='任务不存在') if not job.log_path.exists(): return {'content': '', 'offset': 0, 'next_offset': 0, 'done': job.status not in ('queued', 'running')} data = job.log_path.read_bytes() if offset < 0: offset = 0 if offset > len(data): offset = len(data) chunk = data[offset:].decode('utf-8', errors='replace') return { 'content': chunk, 'offset': offset, 'next_offset': len(data), 'done': job.status not in ('queued', 'running'), 'status': job.status, } @app.get('/api/jobs/{job_id}/stream') async def stream_logs(job_id: str, offset: int = 0): job = JOBS.get(job_id) if not job: raise HTTPException(status_code=404, detail='任务不存在') async def event_gen(): pos = max(0, offset) while True: if job.log_path.exists(): data = job.log_path.read_bytes() if pos < len(data): chunk = data[pos:].decode('utf-8', errors='replace') pos = len(data) payload = json.dumps({'type': 'log', 'content': chunk, 'offset': pos}, ensure_ascii=False) yield f'data: {payload}\n\n' status_payload = json.dumps({ 'type': 'status', 'status': job.status, 'return_code': job.return_code, 'offset': pos, }, ensure_ascii=False) yield f'data: {status_payload}\n\n' if job.status not in ('queued', 'running'): done_payload = json.dumps({'type': 'done', 'status': job.status, 'offset': pos}, ensure_ascii=False) yield f'data: {done_payload}\n\n' break await asyncio.sleep(0.8) return StreamingResponse(event_gen(), media_type='text/event-stream') @app.get('/api/results/overview') async def results_overview(output_dir: Optional[str] = None): root = Path(output_dir) if output_dir else DEFAULT_OUTPUT_DIR return results_scan.scan_output_dir(root) @app.get('/api/results/compare') async def results_compare( models: Optional[str] = None, benchmarks: Optional[str] = None, output_dir: Optional[str] = None, ): root = Path(output_dir) if output_dir else DEFAULT_OUTPUT_DIR folder_list = [x.strip() for x in (models or '').split(',') if x.strip()] or None bench_list = [x.strip() for x in (benchmarks or '').split(',') if x.strip()] or None return results_scan.compare_models(root, folders=folder_list, benchmarks=bench_list) @app.get('/') async def index(): return FileResponse(STATIC_DIR / 'index.html') @app.get('/results') async def results_page(): return FileResponse(STATIC_DIR / 'results.html') app.mount('/static', StaticFiles(directory=str(STATIC_DIR)), name='static') if __name__ == '__main__': import uvicorn host = os.environ.get('WEBUI_HOST', '0.0.0.0') port = int(os.environ.get('WEBUI_PORT', '7860')) uvicorn.run('server:app', host=host, port=port, reload=False)