"""Bridge ML-Alpha-Research-System into web_development API.""" from __future__ import annotations import asyncio import sys from pathlib import Path from typing import Any from fastapi import APIRouter, HTTPException from pydantic import BaseModel, Field WORKSPACE_ROOT = Path(__file__).resolve().parents[4] if str(WORKSPACE_ROOT) not in sys.path: sys.path.insert(0, str(WORKSPACE_ROOT)) router = APIRouter(prefix="/research", tags=["research"]) def _ensure_research_enabled(): from app.services.platform_manager import platform_manager if not platform_manager.qlib_research_enabled: raise HTTPException( status_code=503, detail="Qlib 研究模块未启动。请先 POST /api/platform/services/qlib_research/start", ) class BacktestRequest(BaseModel): strategy: str = "topk_dropout" segment: str = "test" start: str | None = None end: str | None = None class FactorAddRequest(BaseModel): name: str expression: str description: str = "" tags: str = "custom" class OperatorTreeRequest(BaseModel): name: str = "custom_factor" tree: dict[str, Any] = Field(default_factory=dict) @router.get("/health") async def research_health(): from app.services.platform_manager import platform_manager return { "status": "ok" if platform_manager.qlib_research_enabled else "stopped", "workspace": str(WORKSPACE_ROOT), "modules": ["factor_registry", "operator_builder", "strategies", "quantaalpha"], } @router.get("/factors") async def list_factors(enabled_only: bool = True): _ensure_research_enabled() from factor_engine.formula_registry import list_factor_specs loop = asyncio.get_event_loop() specs = await loop.run_in_executor(None, lambda: list_factor_specs(enabled_only=enabled_only)) factors = [ { "name": s.name, "expression": getattr(s, "expression", ""), "description": s.description, "tags": s.tags or [], "category": (s.tags or ["custom"])[0], "enabled": s.enabled, } for s in specs ] return {"factors": factors} @router.post("/factors/{name}/compute") async def compute_factor_api(name: str): _ensure_research_enabled() from factor_engine.formula_registry import compute_factor loop = asyncio.get_event_loop() try: series = await loop.run_in_executor(None, lambda: compute_factor(name, cache=True)) return {"name": name, "count": int(series.notna().sum()), "cached": True} except Exception as exc: raise HTTPException(status_code=400, detail=str(exc)) @router.post("/factors/{name}/backtest") async def backtest_factor_api(name: str, body: BacktestRequest | None = None): _ensure_research_enabled() from factor_engine.single_factor_backtest import backtest_single_factor body = body or BacktestRequest() loop = asyncio.get_event_loop() try: result = await loop.run_in_executor( None, lambda: backtest_single_factor( name, strategy_name=body.strategy, segment=body.segment, start_time=body.start, end_time=body.end, ), ) return result except Exception as exc: raise HTTPException(status_code=400, detail=str(exc)) @router.get("/factors/summary") async def factors_ic_summary(): _ensure_research_enabled() import pandas as pd from config.settings import load_settings path = load_settings().factor_registry_output_dir() / "ic_summary.csv" if not path.exists(): return [] df = pd.read_csv(path) return df.to_dict(orient="records") @router.post("/factors") async def add_factor_api(body: FactorAddRequest): _ensure_research_enabled() from factor_engine.formula_registry import add_factor_to_registry loop = asyncio.get_event_loop() tags = [t.strip() for t in body.tags.split(",") if t.strip()] spec = await loop.run_in_executor( None, lambda: add_factor_to_registry(body.name, body.expression, body.description, tags), ) return {"name": spec.name, "expression": spec.expression} @router.get("/operators") async def list_operators(): _ensure_research_enabled() from factor_engine.operator_builder import get_operator_catalog, list_built_factors catalog = get_operator_catalog() built = [ {"name": s.name, "description": s.description, "enabled": s.enabled} for s in list_built_factors(enabled_only=False) ] operators = [] for field in catalog.get("fields", []): operators.append({"name": field, "category": "field", "description": "行情字段"}) for op in catalog.get("unary", []): operators.append({"name": op, "category": "unary", "description": "一元算子"}) for op in catalog.get("binary", []): operators.append({"name": op, "category": "binary", "description": "二元算子"}) for op in catalog.get("rolling", []): operators.append({"name": op, "category": "rolling", "description": "滚动窗口算子"}) return {"catalog": catalog, "built_factors": built, "operators": operators} @router.post("/operators/build") async def build_operator_tree(body: OperatorTreeRequest): _ensure_research_enabled() from factor_engine.operator_builder import tree_to_qlib_expression loop = asyncio.get_event_loop() try: expr = await loop.run_in_executor(None, lambda: tree_to_qlib_expression(body.tree)) return {"name": body.name, "expression": expr} except Exception as exc: raise HTTPException(status_code=400, detail=str(exc)) @router.post("/operators/{name}/build") async def build_operator_factor(name: str, register: bool = True): _ensure_research_enabled() from factor_engine.operator_builder import build_expression, register_built_factor loop = asyncio.get_event_loop() try: if register: expr = await loop.run_in_executor(None, lambda: register_built_factor(name)) else: expr = await loop.run_in_executor(None, lambda: build_expression(name)) return {"name": name, "expression": expr, "registered": register} except Exception as exc: raise HTTPException(status_code=400, detail=str(exc)) @router.get("/strategies") async def list_research_strategies(): _ensure_research_enabled() from strategies.registry import list_strategies items = list_strategies() strategies = [ { "name": s.get("name"), "description": s.get("description", ""), "class": s.get("class_path", "builtins"), } for s in items ] return {"strategies": strategies} @router.post("/strategies/backtest") async def run_strategy_backtest_api(body: dict[str, Any]): _ensure_research_enabled() from strategies.runner import run_strategy_backtest, save_backtest_result from config.settings import load_settings signal_source = body.get("signal_source") or { "type": "factor_registry", "name": body.get("factor_name", "momentum_5d"), } loop = asyncio.get_event_loop() try: result = await loop.run_in_executor( None, lambda: run_strategy_backtest( strategy_name=body.get("strategy", "topk_dropout"), signal_source=signal_source, start_time=body.get("start"), end_time=body.get("end"), ), ) out = load_settings().output_root / "backtest" / body.get("strategy", "topk_dropout") await loop.run_in_executor(None, lambda: save_backtest_result(result, out)) return { "strategy": result.strategy_name, "signal_stats": result.signal_stats, "risk": result.risk.to_dict() if not result.risk.empty else {}, "output_dir": str(out), } except Exception as exc: raise HTTPException(status_code=400, detail=str(exc)) @router.get("/strategies/suites") async def list_strategy_suites(): _ensure_research_enabled() import yaml from config.settings import load_settings path = load_settings().strategy_config_path with open(path, encoding="utf-8") as f: data = yaml.safe_load(f) return data.get("suites", {})