Files
tick-stock-panel/backend/tests/test_mining_manager.py
T
shy3130 697c27bb02 feat(v0.2): 市场阶段与主线识别 + 因子挖掘全链路 + 数据层完善
- 市场环境: 新增情绪周期6阶段(冰点/启动/主升/高潮/退潮/修复, 连板梯队驱动,
  EMA平滑+2日确认+弱档否决, 平均段长9.7天)与概念/行业主线排名(涨停梯队聚合,
  可配置宽基/风格标签过滤); 市场环境页重构, regime 透明加列, 与5档state并存
- 挖掘: 因子与策略挖掘全链路(API/worker/进程锁/候选库/前端工作台/文档),
  周度调度默认关闭且永不自动发布
- 回测: 财务快照因子(点时口径), 批量回测预计算共享下期收益,
  信号路径矩阵列依赖展开修复(consecutive_limit_ups 缺列报错)
- 数据/性能: enriched 生成与预热治理, 重任务限流, 行情/K线缓存复用, 时区修复
- 测试: 后端全量 914 通过; GUI 黑盒验证截图存证 gui-test-screenshots/
2026-08-16 23:39:07 +08:00

353 lines
11 KiB
Python

from __future__ import annotations
import threading
import time
from collections.abc import Callable
from pathlib import Path
from typing import Any
import pytest
import app.services.mining_manager as mining_manager_module
from app.services.heavy_job_limiter import HeavyJobLimiter
from app.services.mining_manager import MiningJobManager
def _task_factory(kind: str, data_dir: Path, payload: dict[str, Any]) -> dict[str, Any]:
return {
"kind": kind,
"data_dir": str(data_dir),
"payload": payload,
}
def _wait_for_status(
manager: MiningJobManager,
run_id: str,
status: str,
*,
timeout: float = 2.0,
) -> dict[str, Any]:
deadline = time.monotonic() + timeout
while time.monotonic() < deadline:
manifest = manager.store.get(run_id)
assert manifest is not None
if manifest["status"] == status:
return manifest
time.sleep(0.005)
pytest.fail(f"run {run_id} did not reach {status}")
@pytest.fixture
def isolated_limiter(monkeypatch: pytest.MonkeyPatch) -> HeavyJobLimiter:
limiter = HeavyJobLimiter(capacity=2, cancel_poll_interval=0.005)
monkeypatch.setattr(mining_manager_module, "shared_heavy_job_limiter", limiter)
return limiter
@pytest.fixture
def make_manager(
tmp_path: Path,
isolated_limiter: HeavyJobLimiter,
):
managers: list[MiningJobManager] = []
def factory(
runner: Callable[
[dict[str, Any], Callable[[dict[str, Any]], None], threading.Event],
dict[str, Any],
],
) -> MiningJobManager:
manager = MiningJobManager(
tmp_path,
worker_runner=runner,
task_factory=_task_factory,
)
managers.append(manager)
return manager
yield factory
for manager in managers:
manager.shutdown()
assert isolated_limiter.in_use == 0
def test_start_records_states_events_progress_and_worker_payload(
make_manager, tmp_path: Path
) -> None:
progress_recorded = threading.Event()
finish = threading.Event()
tasks: list[dict[str, Any]] = []
progress = {"phase": "screen", "done": 1, "total": 2}
result = {"status": "succeeded", "candidate_count": 3, "elapsed_ms": 12.5}
def runner(task, progress_cb, cancel_event):
tasks.append(task)
progress_cb(progress)
progress_recorded.set()
assert finish.wait(2)
assert not cancel_event.is_set()
return result
manager = make_manager(runner)
request = {"factor_names": ["momentum"], "budget_profile": "balanced"}
created = manager.start(request, {"daily": "v1"}, source="scheduled")
run_id = created["run_id"]
assert created["status"] == "queued"
assert progress_recorded.wait(1)
assert manager.store.read_summary(run_id) == {"progress": progress}
assert [event["type"] for event in manager.store.read_events(run_id)] == [
"queued",
"running",
"progress",
]
assert manager.store.read_events(run_id)[0]["payload"]["source"] == "scheduled"
assert tasks == [
{
"kind": "mining",
"data_dir": str(tmp_path.resolve()),
"payload": {
"run_id": run_id,
"request": request,
"data_fingerprint": {"daily": "v1"},
"source": "scheduled",
},
}
]
finish.set()
terminal = _wait_for_status(manager, run_id, "succeeded")
assert terminal["started_at"] is not None
assert terminal["finished_at"] is not None
assert manager.store.read_summary(run_id) == result
assert [event["type"] for event in manager.store.read_events(run_id)] == [
"queued",
"running",
"progress",
"succeeded",
]
def test_start_accepts_valid_persistent_run_id(make_manager) -> None:
def runner(task, progress_cb, cancel_event):
return {"status": "succeeded"}
manager = make_manager(runner)
created = manager.start(
{"factor_names": ["value"]},
"data-v1",
run_id="weekly_claim_2026_33",
)
assert created["run_id"] == "weekly_claim_2026_33"
terminal = _wait_for_status(manager, created["run_id"], "succeeded")
assert terminal["run_id"] == "weekly_claim_2026_33"
def test_duplicate_persistent_run_id_reuses_existing_without_starting_runner(
make_manager,
) -> None:
runner_called = threading.Event()
def runner(task, progress_cb, cancel_event):
runner_called.set()
return {"status": "succeeded"}
manager = make_manager(runner)
existing = manager.store.create(
{"factor_names": ["value"]},
"data-v1",
run_id="weekly_claim_2026_33",
)
reused = manager.start(
{"factor_names": ["value"]},
"data-v1",
run_id="weekly_claim_2026_33",
)
assert reused == existing
assert not runner_called.wait(0.05)
def test_start_reuses_active_and_success_but_force_creates_new_run(make_manager) -> None:
started = threading.Event()
release = threading.Event()
tasks: list[dict[str, Any]] = []
def runner(task, progress_cb, cancel_event):
tasks.append(task)
started.set()
assert release.wait(2)
return {"status": "succeeded", "candidate_count": 1}
manager = make_manager(runner)
request = {"factor_names": ["value"]}
first = manager.start(request, "data-v1")
assert started.wait(1)
active_reuse = manager.start(request, "data-v1")
assert active_reuse["run_id"] == first["run_id"]
assert len(tasks) == 1
release.set()
_wait_for_status(manager, first["run_id"], "succeeded")
success_reuse = manager.start(request, "data-v1")
assert success_reuse["run_id"] == first["run_id"]
assert len(tasks) == 1
forced = manager.start(request, "data-v1", force=True)
assert forced["run_id"] != first["run_id"]
_wait_for_status(manager, forced["run_id"], "succeeded")
assert len(tasks) == 2
assert len(manager.store.list_runs()) == 2
def test_cancel_while_waiting_for_capacity_never_calls_runner(
make_manager,
isolated_limiter: HeavyJobLimiter,
) -> None:
runner_called = threading.Event()
def runner(task, progress_cb, cancel_event):
runner_called.set()
return {"status": "succeeded"}
assert isolated_limiter.acquire("mining", timeout=0)
try:
manager = make_manager(runner)
created = manager.start({"factor_names": ["quality"]}, "data-v1")
run_id = created["run_id"]
assert manager.store.get(run_id)["status"] == "queued" # type: ignore[index]
cancelling = manager.cancel(run_id)
assert cancelling["status"] == "cancelling"
_wait_for_status(manager, run_id, "cancelled")
assert not runner_called.is_set()
assert [event["type"] for event in manager.store.read_events(run_id)] == [
"queued",
"cancelling",
"cancelled",
]
finally:
isolated_limiter.release("mining")
def test_cancel_running_job_wins_over_worker_success(make_manager) -> None:
runner_started = threading.Event()
def runner(task, progress_cb, cancel_event):
runner_started.set()
assert cancel_event.wait(2)
return {"status": "succeeded", "candidate_count": 9}
manager = make_manager(runner)
created = manager.start({"factor_names": ["growth"]}, "data-v1")
run_id = created["run_id"]
assert runner_started.wait(1)
cancelling = manager.cancel(run_id)
assert cancelling["status"] == "cancelling"
_wait_for_status(manager, run_id, "cancelled")
event_types = [event["type"] for event in manager.store.read_events(run_id)]
assert event_types == ["queued", "running", "cancelling", "cancelled"]
assert manager.store.read_summary(run_id) == {}
def test_runner_exception_marks_failed_and_appends_error_event(make_manager) -> None:
def runner(task, progress_cb, cancel_event):
raise RuntimeError("mining exploded")
manager = make_manager(runner)
created = manager.start({"factor_names": ["size"]}, "data-v1")
run_id = created["run_id"]
failed = _wait_for_status(manager, run_id, "failed")
assert failed["error"] == "mining exploded"
events = manager.store.read_events(run_id)
assert [event["type"] for event in events] == ["queued", "running", "error"]
assert events[-1]["payload"] == {
"status": "failed",
"message": "mining exploded",
}
def test_non_dict_worker_result_is_rejected(make_manager) -> None:
def runner(task, progress_cb, cancel_event):
return ["full", "result"]
manager = make_manager(runner)
created = manager.start({"factor_names": ["liquidity"]}, "data-v1")
failed = _wait_for_status(manager, created["run_id"], "failed")
assert failed["error"] == "mining worker result must be a compact dict"
def test_budget_exhausted_result_uses_distinct_success_status(make_manager) -> None:
result = {
"status": "succeeded_with_budget_exhausted",
"candidate_count": 2,
"budget_exhausted": True,
}
def runner(task, progress_cb, cancel_event):
return result
manager = make_manager(runner)
created = manager.start({"factor_names": ["volatility"]}, "data-v1")
run_id = created["run_id"]
_wait_for_status(manager, run_id, "succeeded_with_budget_exhausted")
assert manager.store.read_summary(run_id) == result
assert manager.store.read_events(run_id)[-1]["type"] == ("succeeded_with_budget_exhausted")
def test_recover_interrupted_delegates_to_store(make_manager) -> None:
def runner(task, progress_cb, cancel_event):
return {"status": "succeeded"}
manager = make_manager(runner)
manager.store.create({}, "v1", run_id="running_before_restart")
manager.store.transition_status("running_before_restart", "running")
manager.store.create({}, "v1", run_id="cancelling_before_restart")
manager.store.transition_status("cancelling_before_restart", "cancelling")
manager.store.create({}, "v1", run_id="queued_before_restart")
assert manager.recover_interrupted() == 3
assert manager.store.get("running_before_restart")["status"] == "interrupted" # type: ignore[index]
assert manager.store.get("cancelling_before_restart")["status"] == "interrupted" # type: ignore[index]
assert manager.store.get("queued_before_restart")["status"] == "interrupted" # type: ignore[index]
def test_shutdown_sets_cancel_uses_bounded_join_and_keeps_history(
make_manager,
monkeypatch: pytest.MonkeyPatch,
) -> None:
runner_started = threading.Event()
release_runner = threading.Event()
def runner(task, progress_cb, cancel_event):
runner_started.set()
assert release_runner.wait(2)
return {"status": "succeeded"}
monkeypatch.setattr(mining_manager_module, "_SHUTDOWN_JOIN_SECONDS", 0.02)
manager = make_manager(runner)
created = manager.start({"factor_names": ["reversal"]}, "data-v1")
run_id = created["run_id"]
assert runner_started.wait(1)
started = time.monotonic()
manager.shutdown()
elapsed = time.monotonic() - started
assert elapsed < 0.2
assert manager.store.get(run_id)["status"] == "cancelling" # type: ignore[index]
release_runner.set()
_wait_for_status(manager, run_id, "cancelled")
assert manager.store.get(run_id) is not None
with pytest.raises(RuntimeError, match="shut down"):
manager.start({"factor_names": ["new"]}, "data-v1")