feat: 增加仅检查 MCP 工具并完善已提交 Alpha 指标展示
Deploy production / deploy (push) Successful in 57s

This commit is contained in:
yuxuanhui
2026-09-11 13:14:27 +08:00
parent f25161d624
commit 7e990b9a69
12 changed files with 311 additions and 63 deletions
+6 -1
View File
@@ -42,18 +42,23 @@ def snapshot_columns(settings, metrics, checks):
check_type = "PENDING"
else:
check_type = "PASS" if "PROD_CORRELATION" in by_name else "PRE_CHECK"
# /check values are freshest; submitted snapshots also expose a scalar in IS.
prod_correlation = number(by_name.get("PROD_CORRELATION", {}).get("value"))
if prod_correlation is None:
prod_correlation = number(metrics.get("prodCorrelation"))
neutralization = settings.get("neutralization")
return {
"check_type": check_type,
"neutralization": neutralization if isinstance(neutralization, str) else None,
"pnl": number(metrics.get("pnl")),
"prod_correlation": prod_correlation,
**{
field: number(by_name.get(name, {}).get("value"))
for field, name in (
("sub_universe_sharpe", "LOW_SUB_UNIVERSE_SHARPE"),
("robust_universe_sharpe", "LOW_ROBUST_UNIVERSE_SHARPE"),
("two_year_sharpe", "LOW_2Y_SHARPE"),
("prod_correlation", "PROD_CORRELATION"),
)
},
}
+4 -2
View File
@@ -21,13 +21,15 @@ from ..research_access.service import ResearchAccess, ResearchError
# Name, schema, business method, required scope, description. No generic arbitrary HTTP tool.
TOOLS = {
"get_submission_check": (c.SelfCorrelationReference, "submission_check_context", "research:read", "读取已导入 Alpha 的表达式、Description、snapshot 和缓存检查结果;不发起检查。先核对或生成三段 Description,再调用 check_submission。"),
"check_submission": (c.SubmissionCheck, "check_submission", "research:refresh", "对单个待提交 Alpha 写回已获用户授权的 Description 并调用平台 GET /check,返回 job_id。须先用 get_submission_check 获取 snapshot;保留本地自相关门槛和冲突保护。通过 get_refresh_job 查进度、get_submission_check 读结果。无论检查结果如何,都不会调用 /submit 或正式提交 Alpha。"),
"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", "查询研究刷新或本地自相关任务的状态、进度与分页错误。"),
"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 按完整输入精确匹配,不推断数学等价。"),
@@ -68,7 +70,7 @@ class MCPResearchServer:
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"}))
openWorldHint=method in {"refresh", "submit", "metadata", "check_self_correlation", "check_submission", "authenticate"}))
for name, (schema, method, scope, description) in TOOLS.items()
if scope in principal.scopes and "research:read" in principal.scopes])
+6
View File
@@ -141,6 +141,12 @@ class SelfCorrelationReference(Contract):
alpha_id: AlphaId
class SubmissionCheck(Contract):
alpha_id: AlphaId
snapshot: str = Field(pattern=r"^[a-f0-9]{64}$")
descriptions: dict[str, str] = Field(min_length=1, max_length=2)
class History(Page):
source: str | None = Field(default=None, max_length=100)
reference: str | None = Field(default=None, max_length=200)
+31 -1
View File
@@ -22,6 +22,8 @@ from ..models import Account, Alpha, BacktestItem, Job, JobItem, ResearchRequest
from ..research.serialization import encode_snapshot
from ..research.workspace_contracts import FieldAvailabilityInput
from ..schemas import JobInput
from ..submission import CheckInput, correlation_allows_check, create_check_job, local_alpha, source
from ..submission import fingerprint as submission_fingerprint
from .contracts import DirectCandidate, History
from .queries import EvidenceQueries, page
@@ -51,6 +53,11 @@ class ResearchAccess:
"confirmation": "调用者须已获本批执行授权;直接提交后返回稳定运行 ID",
"duplicate_policies": ["reject", "rerun"], "permissions": sorted(self.principal.scopes),
"metadata_only": True, "actual_platform_allowance": None,
"submission_check": {
"check_with": "check_submission", "read_with": "get_submission_check",
"job_with": "get_refresh_job", "max_targets": 1,
"description_writeback": True, "production_submission": False,
},
"self_correlation": {
"check_with": "check_self_correlation", "read_with": "get_self_correlation",
"job_with": "get_refresh_job", "max_targets": 100, "source": "local",
@@ -156,7 +163,7 @@ class ResearchAccess:
async def refresh_job(self, args):
job = await self.db.get(Job, args.job_id)
if not job or job.kind not in {"catalog_sync", "field_sync", "pnl_refresh", "self_correlation"}:
if not job or job.kind not in {"catalog_sync", "field_sync", "pnl_refresh", "self_correlation", "submission_check"}:
raise ResearchError("NOT_FOUND", "研究刷新任务不存在")
result = await self.business.get_job_status(args.job_id)
query = select(JobItem).where(JobItem.job_id == job.id, JobItem.error.is_not(None))
@@ -187,6 +194,29 @@ class ResearchAccess:
status = "not_cached" if not data["cached"] else "stale" if data["result"]["stale"] else "available"
return {"alpha_id": args.alpha_id, "source": "local", "status": status, **data}
async def submission_check_context(self, args):
"""Read cached review context and check evidence without upstream requests."""
alpha = await local_alpha(self.db, args.alpha_id)
context = source(alpha.raw)
job = await self.db.scalar(select(Job).where(
Job.kind == "submission_check", Job.payload["alpha_ids"][0].as_string() == args.alpha_id
).order_by(Job.created_at.desc()).limit(1))
return {"alpha_id": args.alpha_id, "snapshot": submission_fingerprint(context),
**context, "descriptions": {key: item["description"] for key, item in context["sections"].items()},
"can_check": alpha.status == "UNSUBMITTED" and await correlation_allows_check(self.db, args.alpha_id),
"checks": alpha.checks, "source": "local_cache", "production_submission": False,
"job_id": job.id if job else None, "job_status": job.status if job else None,
"checked_at": job.checkpoint.get("checked_at") if job else None}
async def check_submission(self, args):
"""Queue the shared reviewed-description/check flow; never submit an Alpha."""
job = await create_check_job(self.db, args.alpha_id, CheckInput(
snapshot=args.snapshot, descriptions=args.descriptions))
self.wake = "jobs"
return {"job_id": job.id, "status": job.status, "alpha_id": args.alpha_id,
"production_submission": False, "job_with": "get_refresh_job",
"read_with": "get_submission_check"}
async def history(self, args):
return await self.evidence.history(args)
+43 -33
View File
@@ -143,6 +143,48 @@ async def require_correlation(db, alpha_id):
raise HTTPException(409, "有本地比较基准时,请先取得有效、样本完整且低于阈值的本地自相关结果")
async def create_check_job(db, alpha_id, body):
"""Queue description writeback and checks only; caller owns commit and wake-up.
HTTP and MCP share locking, conflict checks and active-job deduplication.
No production submission operation is available in this flow.
"""
account = await db.scalar(select(Account).where(Account.id == 1).with_for_update())
if not account.password_encrypted or account.connection_status in ("disconnected", "error"):
raise HTTPException(409, "请先连接 WorldQuant")
alpha = await local_alpha(db, alpha_id)
context = source(alpha.raw)
if alpha.status != "UNSUBMITTED":
raise HTTPException(409, "仅对待提交 Alpha 写回 Description 并检查")
if fingerprint(context) != body.snapshot:
raise HTTPException(409, "Alpha 内容已变化,请重新载入并核对描述")
if set(body.descriptions) != set(context["sections"]):
raise HTTPException(422, "Description 必须匹配 Alpha 的 regular 或 selection/combo 部分")
await require_correlation(db, alpha_id)
texts = {}
for key, draft in body.descriptions.items():
if isinstance(draft, str):
texts[key] = draft
else:
original = context["sections"][key]["description"]
# Preserve existing formatting for older clients sending separate fields.
texts[key] = (
original if parse_description(original) == draft.model_dump() else draft.text()
)
payload = {"alpha_ids": [alpha_id], "expected": context, "descriptions": texts}
for job in (
await db.scalars(select(Job).where(Job.kind == "submission_check", Job.status.in_(ACTIVE)))
).all():
if job.payload.get("alpha_ids") == [alpha_id]:
if job.payload == payload:
return job
raise HTTPException(409, "此 Alpha 已有平台检查任务,请等待完成或取消后再修改描述")
job = Job(id=str(uuid4()), kind="submission_check", payload=payload, total=1)
db.add(job)
await db.flush()
return job
def router(runner, ai):
api = APIRouter(prefix="/api/v1/alphas", tags=["submission-check"], dependencies=[Depends(require_auth)])
generation_lock = asyncio.Lock()
@@ -248,39 +290,7 @@ def router(runner, ai):
@api.post("/{alpha_id}/submission-check", status_code=202, response_model=JobOutput)
async def check(alpha_id: str, body: CheckInput):
async with ai.sessions.begin() as db:
account = await db.scalar(select(Account).where(Account.id == 1).with_for_update())
if not account.password_encrypted or account.connection_status in ("disconnected", "error"):
raise HTTPException(409, "请先连接 WorldQuant")
alpha = await local_alpha(db, alpha_id)
context = source(alpha.raw)
if alpha.status != "UNSUBMITTED":
raise HTTPException(409, "仅对待提交 Alpha 写回 Description 并检查")
if fingerprint(context) != body.snapshot:
raise HTTPException(409, "Alpha 内容已变化,请重新载入并核对描述")
if set(body.descriptions) != set(context["sections"]):
raise HTTPException(422, "Description 必须匹配 Alpha 的 regular 或 selection/combo 部分")
await require_correlation(db, alpha_id)
texts = {}
for key, draft in body.descriptions.items():
if isinstance(draft, str):
texts[key] = draft
else:
original = context["sections"][key]["description"]
# Preserve existing formatting for older clients sending separate fields.
texts[key] = (
original if parse_description(original) == draft.model_dump() else draft.text()
)
payload = {"alpha_ids": [alpha_id], "expected": context, "descriptions": texts}
for job in (
await db.scalars(select(Job).where(Job.kind == "submission_check", Job.status.in_(ACTIVE)))
).all():
if job.payload.get("alpha_ids") == [alpha_id]:
if job.payload == payload:
return job
raise HTTPException(409, "此 Alpha 已有平台检查任务,请等待完成或取消后再修改描述")
job = Job(id=str(uuid4()), kind="submission_check", payload=payload, total=1)
db.add(job)
await db.flush()
job = await create_check_job(db, alpha_id, body)
runner.wake.set()
return job
@@ -0,0 +1,42 @@
"""Backfill production correlation from saved IS metrics without platform calls."""
import math
import sqlalchemy as sa
from alembic import op
revision = "0013"
down_revision = "0012"
branch_labels = None
depends_on = None
def upgrade():
table = sa.table("alphas", sa.column("id", sa.String()),
sa.column("is_metrics", sa.JSON()), sa.column("prod_correlation", sa.Float()))
connection = op.get_bind()
last_id = None
while True:
query = sa.select(table.c.id, table.c.is_metrics).where(table.c.prod_correlation.is_(None)).order_by(table.c.id).limit(500)
if last_id is not None:
query = query.where(table.c.id > last_id)
rows = connection.execute(query).mappings().all()
if not rows:
break
for row in rows:
metrics = row["is_metrics"]
value = metrics.get("prodCorrelation") if isinstance(metrics, dict) else None
if value is None or isinstance(value, bool):
continue
try:
value = float(value)
except (TypeError, ValueError):
continue
if math.isfinite(value):
connection.execute(table.update().where(table.c.id == row["id"]).values(prod_correlation=value))
last_id = rows[-1]["id"]
def downgrade():
# A data correction has no schema to undo; preserve recovered observations.
pass
+1 -1
View File
@@ -162,7 +162,7 @@ async def test_official_sdk_client_and_error_contract(mcp_app):
async with ClientSession(streams[0], streams[1]) as client:
await client.initialize()
listed = await client.list_tools()
assert len(listed.tools) == 15
assert len(listed.tools) == 17
caps = await client.call_tool("get_research_capabilities", {})
assert caps.structured_content["max_candidates"] == 100
result = await client.call_tool("submit_backtests", submission())
+63
View File
@@ -0,0 +1,63 @@
"""Exercise MCP through the durable runner; the fake rejects any submit request."""
import pytest
from fastapi import HTTPException
from app.models import Alpha, SelfCorrelation
from tests.conftest import alpha
from tests.test_mcp import credentials, invoke
from tests.test_mcp import mcp_app as mcp_app
from tests.test_submission import FIELDS, Description, setup
@pytest.mark.parametrize("kind", ["REGULAR", "SUPER"])
@pytest.mark.parametrize("result", ["PASS", "FAIL"])
async def test_mcp_check_never_submits(mcp_app, kind, result, monkeypatch):
from tests import test_submission
checks = [{"name": "PROD_CORRELATION", "result": result}]
monkeypatch.setattr(test_submission, "CHECKS", checks)
sections = ["regular"] if kind == "REGULAR" else ["selection", "combo"]
platform = await setup(mcp_app, alpha(type=kind, **{s: {"code": "rank(close)"} for s in sections}))
platform.pending = 1
principal, _ = await credentials(mcp_app)
caps = await invoke(mcp_app, principal, "get_research_capabilities")
assert caps["submission_check"]["production_submission"] is False
state = await invoke(mcp_app, principal, "get_submission_check", {"alpha_id": "alpha1"})
assert platform.calls == []
args = {"alpha_id": "alpha1", "snapshot": state["snapshot"],
"descriptions": dict.fromkeys(sections, Description(**FIELDS).text())}
started = await invoke(mcp_app, principal, "check_submission", args)
replay = await invoke(mcp_app, principal, "check_submission", args)
assert replay["job_id"] == started["job_id"]
await mcp_app.state.runner.execute(started["job_id"])
job = await invoke(mcp_app, principal, "get_refresh_job", {"job_id": started["job_id"]})
assert job["status"] == "completed", job
data = await invoke(mcp_app, principal, "get_submission_check", {"alpha_id": "alpha1"})
assert data["checks"] == checks and data["checked_at"]
assert data["production_submission"] is False
assert platform.patches == [{s: {"description": args["descriptions"][s]} for s in sections}]
assert platform.calls.count(("GET", "/alphas/alpha1/check")) == 2
assert set(platform.calls) <= {("POST", "/authentication"), ("GET", "/alphas/alpha1"),
("PATCH", "/alphas/alpha1"), ("GET", "/alphas/alpha1/check")}
async with mcp_app.state.sessions() as db:
assert (await db.get(Alpha, "alpha1")).status == "UNSUBMITTED"
async def test_mcp_check_scope_and_guards(mcp_app):
platform = await setup(mcp_app)
reader, _ = await credentials(mcp_app, {"research:read"})
state = await invoke(mcp_app, reader, "get_submission_check", {"alpha_id": "alpha1"})
args = {"alpha_id": "alpha1", "snapshot": state["snapshot"], "descriptions": {"regular": "reviewed"}}
with pytest.raises(HTTPException) as exc:
await mcp_app.state.mcp.invoke(reader, "check_submission", args)
assert exc.value.status_code == 403
principal, _ = await credentials(mcp_app)
for extra in [{"snapshot": "0" * 64}, {"descriptions": {"combo": "wrong"}}, {"submit": True}]:
response = await mcp_app.state.mcp.invoke(principal, "check_submission", args | extra)
assert response.is_error
async with mcp_app.state.sessions.begin() as db:
(await db.get(SelfCorrelation, "alpha1")).result = {"status": "high"}
response = await mcp_app.state.mcp.invoke(principal, "check_submission", args)
assert response.is_error
assert platform.calls == []
+56
View File
@@ -0,0 +1,56 @@
import pytest
from app.alphas import snapshot_columns, upsert_alpha
from tests.conftest import alpha
@pytest.mark.parametrize("value", [0.6988, 0.0])
async def test_submitted_metrics_map_without_check_history(app, logged_in, value):
raw = alpha(status="ACTIVE", **{"is": {"prodCorrelation": value, "selfCorrelation": 0.0,
"checks": [{"name": "LOW_2Y_SHARPE", "result": "PASS", "value": 2.1}]},
"os": {"checks": [{"name": "PROD_CORRELATION", "result": "PENDING"}]}})
async with app.state.sessions.begin() as db:
await upsert_alpha(db, raw)
detail = (await logged_in.get('/api/v1/alphas/alpha1')).json()
assert detail['prod_correlation'] == value
assert detail['is_metrics']['selfCorrelation'] == 0.0
assert detail['os_metrics']['checks'][0]['result'] == 'PENDING'
assert detail['two_year_sharpe'] == 2.1
assert detail['check_type'] != 'PASS'
def test_check_value_takes_precedence_over_cached_metric():
assert snapshot_columns({}, {'prodCorrelation': 0.4}, [
{'name': 'PROD_CORRELATION', 'value': 0.0, 'result': 'PASS'}
])['prod_correlation'] == 0.0
def test_migration_recovers_saved_values_without_overwriting_existing():
import importlib.util
from pathlib import Path
from types import SimpleNamespace
import sqlalchemy as sa
path = Path(__file__).parents[1] / 'migrations/versions/0013_submitted_correlation.py'
spec = importlib.util.spec_from_file_location('migration_0013', path)
migration = importlib.util.module_from_spec(spec)
spec.loader.exec_module(migration)
engine = sa.create_engine('sqlite://')
metadata = sa.MetaData()
table = sa.Table('alphas', metadata, sa.Column('id', sa.String, primary_key=True),
sa.Column('is_metrics', sa.JSON), sa.Column('prod_correlation', sa.Float))
metadata.create_all(engine)
with engine.begin() as db:
db.execute(table.insert(), [
{'id': 'zero', 'is_metrics': {'prodCorrelation': 0.0}, 'prod_correlation': None},
{'id': 'saved', 'is_metrics': {'prodCorrelation': 0.6988}, 'prod_correlation': None},
{'id': 'existing', 'is_metrics': {'prodCorrelation': 0.4}, 'prod_correlation': 0.2},
{'id': 'missing', 'is_metrics': {}, 'prod_correlation': None},
])
migration.op = SimpleNamespace(get_bind=lambda: db)
migration.upgrade()
migration.upgrade()
assert dict(db.execute(sa.select(table.c.id, table.c.prod_correlation)).all()) == {
'zero': 0.0, 'saved': 0.6988, 'existing': 0.2, 'missing': None}
engine.dispose()