"""D2 选股 Job 化测试:kind=selection 异步执行(executor + API 提交/轮询)。""" from __future__ import annotations from datetime import date, datetime import pytest from app.api import deps from app.application.services import job_executor as je from app.domain.entities.market import Stock from app.domain.entities.research import JobRecord, JobStatus from app.domain.entities.selection import SelectionQuery, SelectionResult from app.infrastructure.persistence.sqlalchemy.base import Base from app.infrastructure.persistence.sqlalchemy.repositories.jobs_impl import ( SqlAlchemyExperimentRepository, SqlAlchemyJobRepository, ) from app.infrastructure.persistence.sqlalchemy.repositories.market_impl import ( SqlAlchemyDailyBarRepository, SqlAlchemyStockRepository, ) from app.main import app from fastapi.testclient import TestClient from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker from conftest_quant import bars_dataframe_to_daily_bars, synthetic_daily _SYMS = ["600000.SH", "600001.SH", "600002.SH", "600003.SH", "600004.SH"] def _query(symbols=_SYMS) -> SelectionQuery: return SelectionQuery( universe={"exclude_st": False, "min_listing_days": 0, "symbols": list(symbols)}, factors=[{"name": "momentum_60", "weight": 1.0}], top_n=3, as_of=date(2024, 12, 31), ) @pytest.fixture() def engine_session(tmp_path): engine = create_engine(f"sqlite:///{tmp_path / 'db.db'}", future=True) Base.metadata.create_all(engine) Session = sessionmaker(bind=engine, expire_on_commit=False) df = synthetic_daily({s: 0.006 - 0.0015 * i for i, s in enumerate(_SYMS)}, n=320) with Session() as session: SqlAlchemyStockRepository(session).upsert_many( [Stock(symbol=s, name=f"测试{i}", list_date=date(1999, 1, 1)) for i, s in enumerate(_SYMS)] ) SqlAlchemyDailyBarRepository(session).upsert_many(bars_dataframe_to_daily_bars(df)) session.commit() return Session class TestSelectionJobExecutor: def test_execute_and_archive(self, engine_session) -> None: Session = engine_session with Session() as session: SqlAlchemyJobRepository(session).create( JobRecord(id="JOB-SEL-1", kind="selection", spec_json=_query().model_dump_json(), status=JobStatus.QUEUED, created_at=datetime.now()) ) session.commit() je.execute_job( "JOB-SEL-1", session_factory=Session, job_repo_factory=lambda s: SqlAlchemyJobRepository(s), experiment_repo_factory=lambda s: SqlAlchemyExperimentRepository(s), stock_repo_factory=lambda s: SqlAlchemyStockRepository(s), daily_repo_factory=lambda s: SqlAlchemyDailyBarRepository(s), engine=None, ) with Session() as session: job = SqlAlchemyJobRepository(session).get("JOB-SEL-1") exp = SqlAlchemyExperimentRepository(session).get(job.experiment_id or "") assert job.status == JobStatus.SUCCESS result = SelectionResult.model_validate_json(job.result_json or "{}") assert len(result.candidates) == 3 assert exp is not None and exp.kind == "selection" assert "as_of" in (exp.summary_text or "") and "选出 3" in (exp.summary_text or "") class TestSelectionJobApi: @pytest.fixture() def client(self, tmp_path, monkeypatch): from app.infrastructure.persistence.sqlalchemy import session as sess_mod engine = create_engine(f"sqlite:///{tmp_path / 'api.db'}", future=True) Base.metadata.create_all(engine) Session = sessionmaker(bind=engine, expire_on_commit=False) monkeypatch.setattr(sess_mod, "SessionLocal", Session) df = synthetic_daily({s: 0.006 - 0.0015 * i for i, s in enumerate(_SYMS)}, n=320) with Session() as session: SqlAlchemyStockRepository(session).upsert_many( [Stock(symbol=s, name=f"测试{i}", list_date=date(1999, 1, 1)) for i, s in enumerate(_SYMS)] ) SqlAlchemyDailyBarRepository(session).upsert_many(bars_dataframe_to_daily_bars(df)) session.commit() def _session_override(): with Session() as s: yield s app.dependency_overrides[deps.get_session] = _session_override with TestClient(app) as c: yield c app.dependency_overrides.clear() def test_submit_poll(self, client) -> None: body = _query().model_dump(mode="json") resp = client.post("/api/selections/jobs", json=body) assert resp.status_code == 200 job_id = resp.json()["job_id"] state = None for _ in range(40): state = client.get(f"/api/jobs/{job_id}").json() if state["status"] in ("success", "failed", "cancelled"): break assert state["status"] == "success" result = state["result"] assert result is not None and len(result["candidates"]) == 3 assert result["statistics"]["selected"] == 3 # 实验归档存在(selection 类型) exps = client.get("/api/experiments").json() assert any(e["kind"] == "selection" for e in exps)