132 lines
9.7 KiB
Python
132 lines
9.7 KiB
Python
"""MCP transport over shared research operations, with minimal durable audit evidence."""
|
||
|
||
import asyncio
|
||
import json
|
||
import time
|
||
from uuid import uuid4
|
||
|
||
import anyio
|
||
from fastapi import HTTPException
|
||
from mcp import types
|
||
from mcp.server.lowlevel import Server
|
||
from mcp.server.transport_security import TransportSecuritySettings
|
||
from pydantic import ValidationError
|
||
|
||
from ..alphas import sanitize
|
||
from ..backtests.contracts import fingerprint
|
||
from ..models import MCPAudit, now
|
||
from ..research.serialization import encode_snapshot
|
||
from ..research_access import contracts as c
|
||
from ..research_access.service import ResearchAccess, ResearchError
|
||
|
||
# Name, schema, business method, required scope, description. No generic arbitrary HTTP tool.
|
||
TOOLS = {
|
||
"get_worldquant_connection": (c.ConnectionReference, "connection", "research:read", "读取 WorldQuant 连接状态及可选认证 job_id 的进度,不发起认证;人工验证在网页完成。"),
|
||
"authenticate_worldquant": (c.Authentication, "authenticate", "research:refresh", "使用服务端已保存凭据连接或重新认证 WorldQuant,返回 job_id;action=connect(默认)或人工验证后 verify。用 get_worldquant_connection 查询,不接收密码,不修改账户配置。"),
|
||
"get_research_capabilities": (c.Empty, "capabilities", "research:read", "读取直接研究能力、完整设置 schema 和调度阻塞,不代表平台剩余额度。"),
|
||
"search_catalog": (c.CatalogSearch, "catalog", "research:read", "分页查询指定范围的数据集或字段元数据;无缓存不等于无数据,不隐式刷新。"),
|
||
"get_research_metadata": (c.Metadata, "metadata", "research:read", "读取范围、设置快照、算子定义或字段可用性;未知不认定通过。"),
|
||
"refresh_research_data": (c.Refresh, "refresh", "research:refresh", "显式刷新目录、算子、设置、字段可用性或 PnL;不会创建模拟。任务返回 job_id。"),
|
||
"get_refresh_job": (c.JobReference, "refresh_job", "research:read", "查询研究刷新或本地自相关任务的状态、进度与分页错误。"),
|
||
"check_self_correlation": (c.SelfCorrelationCheck, "check_self_correlation", "research:refresh", "对 1–100 个已导入 Alpha 发起本地自相关检查,返回 job_id。与本地同地区已提交 Alpha 比较,排除自身;优先用缓存,缺失 PnL 自动补取。需先同步已提交 Alpha;不调用平台提交检查。用 get_refresh_job 查进度、get_self_correlation 读结果。"),
|
||
"get_self_correlation": (c.SelfCorrelationReference, "self_correlation", "research:read", "读取指定 Alpha 最新的本地自相关缓存,包括最大相关系数、样本覆盖与 stale 状态;无缓存不自动检查,结果不等同于平台提交资格。"),
|
||
"search_backtests": (c.History, "history", "research:read", "分页查历史候选与固定设置;candidates 按完整输入精确匹配,不推断数学等价。"),
|
||
"submit_backtests": (c.Submit, "submit", "backtests:execute", "执行用户已授权的固定批次,自动留痕并立即返回运行 ID。每项必须完整设置;重复默认拒绝,rerun 明确重跑。不需要研究资产。"),
|
||
"get_backtest": (c.RunReference, "run", "research:read", "读取真实运行进度、提交数量和可选增量事件;受理不等于成功。"),
|
||
"get_backtest_results": (c.Results, "results", "research:read", "分页读取固定快照指标、全部非通过检查及三层状态;缺失指标不补零。"),
|
||
"get_backtest_artifact": (c.Artifact, "artifact", "research:read", "分页读取候选脱敏快照的顶层键值或独立采集的 PnL;缺缓存不自动刷新。"),
|
||
"control_backtest": (c.Control, "control", "backtests:control", "对已授权运行暂停、继续、停止或恢复采集;不远程取消、不重提未知模拟。需要版本和幂等键。"),
|
||
}
|
||
|
||
|
||
def tool_result(data, error=False):
|
||
data = encode_snapshot(data)
|
||
return types.CallToolResult(content=[types.TextContent(type="text", text=json.dumps(data, ensure_ascii=False))],
|
||
structuredContent=data, isError=error)
|
||
|
||
|
||
class MCPResearchServer:
|
||
def __init__(self, sessions, runner, settings):
|
||
self.sessions, self.runner, self.settings = sessions, runner, settings
|
||
# The existing deployment has one owner; this also gives SQLite test transactions a fair queue.
|
||
self.mutation_lock = asyncio.Lock()
|
||
self.server = Server("wq-alpha-research", version="1.0.0", on_list_tools=self.list_tools,
|
||
on_call_tool=self.call_tool,
|
||
instructions="自由探索,直接固定候选回测,无需先建研究资产。工具不安排定时研究;结果按运行 ID 查询。")
|
||
from urllib.parse import urlsplit
|
||
|
||
host = urlsplit(settings.public_origin).netloc
|
||
self.app = self.server.streamable_http_app(
|
||
streamable_http_path="/", stateless_http=True, json_response=True,
|
||
transport_security=TransportSecuritySettings(enable_dns_rebinding_protection=True,
|
||
allowed_hosts=[host], allowed_origins=[settings.public_origin.rstrip("/")]),
|
||
)
|
||
|
||
async def list_tools(self, ctx, params):
|
||
principal = ctx.request.state.mcp_principal
|
||
return types.ListToolsResult(tools=[types.Tool(name=name, description=description,
|
||
inputSchema=schema.model_json_schema(), annotations=types.ToolAnnotations(
|
||
readOnlyHint=scope == "research:read", destructiveHint=method == "control",
|
||
idempotentHint=method in {"submit", "control"} or scope == "research:read",
|
||
openWorldHint=method in {"refresh", "submit", "metadata", "check_self_correlation", "authenticate"}))
|
||
for name, (schema, method, scope, description) in TOOLS.items()
|
||
if scope in principal.scopes and "research:read" in principal.scopes])
|
||
|
||
async def call_tool(self, ctx, params):
|
||
principal = ctx.request.state.mcp_principal
|
||
return await self.invoke(principal, params.name, params.arguments or {}, str(ctx.request_id or uuid4()))
|
||
|
||
async def invoke(self, principal, name, arguments, request_id=None):
|
||
"""Invoke with a server-authenticated principal; atomic success audit and post-commit wake."""
|
||
started = time.monotonic()
|
||
request_id = request_id or str(uuid4())
|
||
entry = TOOLS.get(name)
|
||
if not entry:
|
||
return tool_result({"error": ResearchError("UNKNOWN_TOOL", "工具不存在").data}, True)
|
||
schema, method, scope, _ = entry
|
||
if "research:read" not in principal.scopes or scope not in principal.scopes:
|
||
raise HTTPException(403, "MCP 令牌缺少所需权限")
|
||
digest = fingerprint(arguments)
|
||
async with self.mutation_lock:
|
||
# Disconnect does not roll back an already accepted operation or lose its wake-up.
|
||
with anyio.CancelScope(shield=True):
|
||
async with self.sessions.begin() as db:
|
||
access = ResearchAccess(db, principal, self.runner.client, self.settings.public_origin)
|
||
code, error = "OK", False
|
||
try:
|
||
async with db.begin_nested():
|
||
args = schema.model_validate(arguments)
|
||
async with asyncio.timeout(30 if method in {"refresh", "metadata"} else None):
|
||
data = encode_snapshot(await getattr(access, method)(args))
|
||
data.setdefault("_meta", {"schema_version": 1, "observed_at": now().isoformat(),
|
||
"nulls": "null 表示来源未提供,不等于零", "source": "system"})
|
||
except ValidationError as exc:
|
||
code, error = "INVALID_INPUT", True
|
||
data = {"error": ResearchError(code, "; ".join(
|
||
f"{'.'.join(map(str, e['loc']))}: {e['msg']}" for e in exc.errors())).data}
|
||
except TimeoutError:
|
||
code, error = "UPSTREAM_TIMEOUT", True
|
||
data = {"error": ResearchError(code, "元数据读取或刷新超时,未发布新快照", retryable=True).data}
|
||
except ResearchError as exc:
|
||
code, error, data = exc.data["code"], True, {"error": exc.data}
|
||
except HTTPException as exc:
|
||
code = {404: "NOT_FOUND", 409: "CONFLICT", 422: "INVALID_INPUT", 429: "RATE_LIMITED", 502: "UPSTREAM_ERROR"}.get(exc.status_code, "REQUEST_FAILED")
|
||
error = True
|
||
data = {"error": ResearchError(code, str(sanitize(exc.detail)),
|
||
retryable=exc.status_code in {429, 502, 503},
|
||
retry_after=(exc.headers or {}).get("Retry-After")).data}
|
||
except Exception:
|
||
# Never expose SQL parameters, exception reprs or credentials in unexpected errors.
|
||
code, error = "INTERNAL_ERROR", True
|
||
data = {"error": ResearchError(code, "研究操作失败;可使用原幂等键重试或查询历史", retryable=True).data}
|
||
db.add(MCPAudit(id=str(uuid4()), token_id=principal.token_id, tool=name,
|
||
request_id=fingerprint({"request_id": request_id}), input_digest=digest,
|
||
business_id=data.get("backtest_run_id", data.get("job_id")),
|
||
result_code=code, elapsed_ms=int((time.monotonic()-started)*1000)))
|
||
if not error:
|
||
if access.wake == "backtests":
|
||
self.runner.backtests.wake.set()
|
||
elif access.wake == "jobs":
|
||
self.runner.wake.set()
|
||
return tool_result(data, error)
|