字段库(本次新增的表与接口): - `condition_field` 表 + `/api/condition-fields`:中文名/说明可编辑、可停用; `kind`/单位阶梯/`base_unit` 由代码注册表收敛(改类型 422,伪字段 422, 越界单位 422),停用的字段不再进条件下拉,但既有策略仍按名字解析。 - 说明书里的数值条件按字段注册表补**基准单位**后缀(字段间比较不加,不猜单位)。 因子参数化(键即身份,冻结口径): - 模板 + 参数注册表(`quant/factors.py`):`ParamSpec`(类型/范围/枚举/默认值/说明)+ `FactorTemplate`(公式/依赖列/参数);规范键把**全部**参数写进名字,如 `momentum(window=90,direction=lower_is_better)`,所以改参数 = 新建一个身份, 旧因子/既有策略/已归档实验都不变义;`momentum(window=90)`(缺参数)明确拒绝 —— 缺项要靠模板默认值补齐,而默认值是可改的代码细节,一旦改动会追溯性改义。 - 参数只在受控范围内取值(窗口 2~500、方向二选一),越界/未知模板/多给参数一律 422 并列出允许范围,不静默截断、不悄悄取默认值;内置实例的启用开关由代码决定(422)。 - `/api/factors` 暴露 `template`/`params`/`param_specs`/`label`/`source`/`enabled`/ `resolvable`;新增 `/api/factors/templates`、`POST /api/factors`、`PATCH /api/factors`; `get_factor = resolve_factor` 兼容全部旧调用点,参数化键也是一等条件字段。 - 迁移链:c5d6(存量策略陈旧说明重算)→ d6e7(condition_field)→ a7c1 (factor_definition.enabled + name varchar(128))。 测试:新增 test_condition_fields.py / test_factor_params.py;全量 pytest 500 passed。
157 lines
5.9 KiB
Python
157 lines
5.9 KiB
Python
"""异步研究 Job API(Phase 4):提交 / 查询 / SSE 进度。
|
||
|
||
POST /api/jobs 创建 Job(BackgroundTasks 后台执行),立即返回 job_id
|
||
GET /api/jobs/{id} 状态 + 结果(成功时内嵌 result)
|
||
GET /api/jobs/{id}/events SSE 进度(queued→running→success|failed)
|
||
|
||
执行模式见 config job.mode:subprocess 时研究任务在独立子进程跑(内存隔离),
|
||
API worker 不被重任务拖垮(内存优化专项)。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import json
|
||
from datetime import datetime
|
||
from typing import Annotated
|
||
|
||
from fastapi import APIRouter, BackgroundTasks, HTTPException, Query
|
||
from fastapi.responses import StreamingResponse
|
||
|
||
from app.api.deps import DbSession, ExperimentRepoDep, JobRepoDep
|
||
from app.application.services.job_executor import new_id, run_job_background, terminate_active
|
||
from app.domain.entities.research import (
|
||
BacktestResult,
|
||
FactorTestReport,
|
||
JobRecord,
|
||
JobStatus,
|
||
ResearchSpec,
|
||
)
|
||
from app.infrastructure.persistence.sqlalchemy.repositories.jobs_impl import (
|
||
SqlAlchemyJobRepository,
|
||
)
|
||
from app.infrastructure.persistence.sqlalchemy.session import SessionLocal
|
||
|
||
router = APIRouter(prefix="/jobs", tags=["jobs"])
|
||
|
||
|
||
def _decode_result(kind: str, result_json: str | None):
|
||
from app.domain.entities.selection import SelectionResult
|
||
|
||
if result_json is None:
|
||
return None
|
||
if kind == "selection":
|
||
return SelectionResult.model_validate_json(result_json)
|
||
# combo 回测的归档 kind 记为 "backtest",但 Job.kind 仍是 "combo" ——
|
||
# 其结果同样是 BacktestResult,按 backtest 解码(否则会被当成因子测试而校验失败)。
|
||
model = BacktestResult if kind in ("backtest", "combo") else FactorTestReport
|
||
return model.model_validate_json(result_json)
|
||
|
||
|
||
def _job_view(job: JobRecord, experiment_repo=None) -> dict:
|
||
"""Job 视图:`result` 契约不变(成功时内嵌**完整**结果)。
|
||
|
||
结果来源(2026-09 起完整结果只在 experiment 存一份,job.result_json 不再重复写):
|
||
1. `job.experiment_id` 有值且能读到归档 → 解码 experiment.result_json;
|
||
2. 否则回退解码 `job.result_json`(老记录 / 归档被删除前的历史数据);
|
||
3. 归档被删除且 job 侧无副本 → `result=None`,并给出
|
||
`result_unavailable_reason` 如实说明原因(AGENT §7:不静默给空结果)。
|
||
"""
|
||
view = job.model_dump()
|
||
view.pop("result_json", None)
|
||
view["spec"] = json.loads(job.spec_json)
|
||
|
||
result = None
|
||
source = None
|
||
if job.experiment_id and experiment_repo is not None:
|
||
exp = experiment_repo.get(job.experiment_id)
|
||
if exp is not None:
|
||
result = _decode_result(job.kind, exp.result_json)
|
||
source = "experiment"
|
||
else:
|
||
view["result_unavailable_reason"] = (
|
||
f"归档 {job.experiment_id} 已不存在(可能已被删除);"
|
||
"完整结果仅存于归档,Job 记录本身不再保存结果副本"
|
||
)
|
||
if result is None and job.result_json:
|
||
result = _decode_result(job.kind, job.result_json)
|
||
source = "job"
|
||
view["result"] = result
|
||
view["result_source"] = source
|
||
return view
|
||
|
||
|
||
@router.post("", summary="创建异步研究 Job")
|
||
def create_job(
|
||
spec: ResearchSpec,
|
||
background: BackgroundTasks,
|
||
session: DbSession,
|
||
job_repo: JobRepoDep,
|
||
) -> dict:
|
||
job = JobRecord(
|
||
id=new_id("JOB"),
|
||
kind=spec.type,
|
||
spec_json=spec.model_dump_json(),
|
||
status=JobStatus.QUEUED,
|
||
created_at=datetime.now(),
|
||
)
|
||
job_repo.create(job)
|
||
session.commit()
|
||
background.add_task(run_job_background, job.id)
|
||
return {"job_id": job.id, "status": job.status}
|
||
|
||
|
||
@router.get("", summary="Job 列表")
|
||
def list_jobs(
|
||
job_repo: JobRepoDep,
|
||
experiment_repo: ExperimentRepoDep,
|
||
kind: Annotated[str | None, Query(description="按类型过滤(backtest/factor_test)")] = None,
|
||
limit: Annotated[int, Query(ge=1, le=200)] = 20,
|
||
) -> list[dict]:
|
||
return [
|
||
_job_view(j, experiment_repo) for j in job_repo.list_recent(kind=kind, limit=limit)
|
||
]
|
||
|
||
|
||
@router.post("/{job_id}/cancel", summary="取消 Job(queued/running)")
|
||
def cancel_job(job_id: str, session: DbSession, job_repo: JobRepoDep) -> dict:
|
||
job = job_repo.get(job_id)
|
||
if job is None:
|
||
raise HTTPException(status_code=404, detail=f"Job {job_id} 不存在")
|
||
if job.status not in (JobStatus.QUEUED, JobStatus.RUNNING):
|
||
return {"job_id": job_id, "status": job.status, "cancelled": False}
|
||
job.status = JobStatus.CANCELLED
|
||
job.stage = None
|
||
job_repo.update(job)
|
||
session.commit()
|
||
terminate_active(job_id) # 终止研究子进程(若有);父进程兜底已跳过 CANCELLED
|
||
return {"job_id": job_id, "status": JobStatus.CANCELLED, "cancelled": True}
|
||
|
||
|
||
@router.get("/{job_id}", summary="查询 Job 状态与结果")
|
||
def get_job(job_id: str, job_repo: JobRepoDep, experiment_repo: ExperimentRepoDep) -> dict:
|
||
job = job_repo.get(job_id)
|
||
if job is None:
|
||
raise HTTPException(status_code=404, detail=f"Job {job_id} 不存在")
|
||
return _job_view(job, experiment_repo)
|
||
|
||
|
||
@router.get("/{job_id}/events", summary="Job 进度 SSE")
|
||
async def job_events(job_id: str) -> StreamingResponse:
|
||
async def gen():
|
||
while True:
|
||
with SessionLocal() as session:
|
||
job = SqlAlchemyJobRepository(session).get(job_id)
|
||
if job is None:
|
||
yield "event: error\ndata: job not found\n\n"
|
||
return
|
||
payload = json.dumps(
|
||
{"job_id": job.id, "status": job.status, "stage": job.stage}, ensure_ascii=False
|
||
)
|
||
yield f"data: {payload}\n\n"
|
||
if job.status in (JobStatus.SUCCESS, JobStatus.FAILED, JobStatus.CANCELLED):
|
||
return
|
||
await asyncio.sleep(0.4)
|
||
|
||
return StreamingResponse(gen(), media_type="text/event-stream")
|