evalstone/webui/server.py
sora 541b0b477d Improve EvalStone web UI paths, custom suites, and results charts.
Move required path fields up front, persist custom benchmark suites, enlarge readable typography, show full suite benchmark names, add capability-domain trend charts with local Chart.js, and polish ranking name display.
2026-07-29 07:38:05 +00:00

643 lines
21 KiB
Python

#!/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)
CUSTOM_SUITES_PATH = DATA_DIR / 'custom_suites.json'
app = FastAPI(title='EvalStone Launch Panel', version='1.0.0')
def _load_custom_suites() -> Dict[str, List[str]]:
if not CUSTOM_SUITES_PATH.exists():
return {}
try:
data = json.loads(CUSTOM_SUITES_PATH.read_text(encoding='utf-8'))
if not isinstance(data, dict):
return {}
out: Dict[str, List[str]] = {}
for name, items in data.items():
if not isinstance(name, str) or not name.strip():
continue
if not isinstance(items, list):
continue
cleaned = [str(x).strip() for x in items if str(x).strip()]
if cleaned:
out[name.strip()] = cleaned
return out
except Exception:
return {}
def _save_custom_suites(suites: Dict[str, List[str]]) -> None:
DATA_DIR.mkdir(parents=True, exist_ok=True)
CUSTOM_SUITES_PATH.write_text(
json.dumps(suites, ensure_ascii=False, indent=2),
encoding='utf-8',
)
class CustomSuiteRequest(BaseModel):
name: str = Field(..., min_length=1)
datasets: List[str] = Field(..., min_length=1)
# ---------------------------------------------------------------------------
# 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,
'custom_suites': _load_custom_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/custom-suites')
async def list_custom_suites():
return {'suites': _load_custom_suites()}
@app.post('/api/custom-suites')
async def upsert_custom_suite(req: CustomSuiteRequest):
name = req.name.strip()
if not name:
raise HTTPException(status_code=400, detail='组合名称不能为空')
if name in run_module.SUITES:
raise HTTPException(status_code=400, detail=f'名称与内置 suite 冲突: {name}')
datasets = []
seen = set()
for item in req.datasets:
d = str(item).strip()
if not d or d in seen:
continue
seen.add(d)
datasets.append(d)
if not datasets:
raise HTTPException(status_code=400, detail='请至少选择一个 benchmark')
suites = _load_custom_suites()
suites[name] = datasets
_save_custom_suites(suites)
return {'ok': True, 'name': name, 'datasets': datasets, 'suites': suites}
@app.delete('/api/custom-suites/{name}')
async def delete_custom_suite(name: str):
suites = _load_custom_suites()
if name not in suites:
raise HTTPException(status_code=404, detail='自定义组合不存在')
suites.pop(name, None)
_save_custom_suites(suites)
return {'ok': True, 'suites': suites}
@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)