Files
Simon 861a4051ca feat(job): D2 全市场选股 Job 化(kind=selection 异步)
- Job kind=selection:executor 分支用 SelectionService(spec=SelectionQuery),结果
  落 result_json、自动归档 Experiment(summary:as_of + 选出 N);/api/jobs 与
  /api/experiments 的 result 解码支持 SelectionResult
- POST /api/selections/jobs:提交 SelectionQuery 为异步 Job(BackgroundTasks/subprocess)
- Web 选股页新增「异步(全市场)」按钮:提交 Job → waitJob 轮询结果(解决同步 60s+)
- tests/test_selection_job.py:executor 执行归档(SelectionResult/Experiment kind)、
  API 提交→轮询→结果与实验列表;全量 pytest + tsc 通过
2026-09-09 07:43:35 +08:00

94 lines
3.0 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""选股 API(M6.3):提交/查询选股,结果落库可复现。
POST /api/selections 同步执行一次选股并落库 → {selection_id, result}
GET /api/selections/{id} 读回某次选股完整结果
GET /api/selections 历史选股元数据(可过滤 as_of/method)
同步执行:单日全市场因子评分/条件计算量轻(秒级);未来若超时再迁 Job。
"""
from __future__ import annotations
from datetime import date, datetime
from typing import Annotated
from fastapi import APIRouter, BackgroundTasks, HTTPException, Query
from pydantic import BaseModel
from app.api.deps import (
DbSession,
JobRepoDep,
SelectionRepoDep,
SelectionServiceDep,
)
from app.application.services.job_executor import new_id, run_job_background
from app.domain.entities.research import JobRecord, JobStatus
from app.domain.entities.selection import SelectionMeta, SelectionQuery, SelectionResult
router = APIRouter(prefix="/selections", tags=["selections"])
class SelectionRun(BaseModel):
selection_id: str
result: SelectionResult
@router.post("/jobs", summary="提交全市场/长任务选股为异步 Job")
def submit_selection_job(
query: SelectionQuery,
background: BackgroundTasks,
session: DbSession,
job_repo: JobRepoDep,
) -> dict:
job = JobRecord(
id=new_id("JOB"),
kind="selection",
spec_json=query.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.post("", response_model=SelectionRun, summary="执行一次选股(同步)并落库")
def run_selection(
query: SelectionQuery,
service: SelectionServiceDep,
selection_repo: SelectionRepoDep,
session: DbSession,
) -> SelectionRun:
result = service.select(query)
selection_id = new_id("SEL")
selection_repo.save(selection_id, result)
session.commit()
return SelectionRun(selection_id=selection_id, result=result)
@router.get("/{selection_id}", response_model=SelectionResult, summary="读回一次选股结果")
def get_selection(
selection_id: str,
selection_repo: SelectionRepoDep,
) -> SelectionResult:
result = selection_repo.get(selection_id)
if result is None:
raise HTTPException(status_code=404, detail=f"选股记录 {selection_id} 不存在")
return result
_AsOfQuery = Annotated[date | None, Query(description="按选股时点过滤")]
_MethodQuery = Annotated[str | None, Query(pattern="^(score|condition)$")]
_LimitQuery = Annotated[int, Query(ge=1, le=200)]
@router.get("", response_model=list[SelectionMeta], summary="历史选股元数据列表")
def list_selections(
selection_repo: SelectionRepoDep,
as_of: _AsOfQuery = None,
method: _MethodQuery = None,
limit: _LimitQuery = 20,
) -> list[SelectionMeta]:
return selection_repo.list_recent(as_of=as_of, method=method, limit=limit)