| """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", {}) |
|
|