diff --git a/.scratch/submitted-pnl/issues/01-backfill.md b/.scratch/submitted-pnl/issues/01-backfill.md new file mode 100644 index 0000000..d996d81 --- /dev/null +++ b/.scratch/submitted-pnl/issues/01-backfill.md @@ -0,0 +1,12 @@ +# 实现已提交 Alpha 检查 PnL + +Type: task +Status: ready-for-agent + +实现 spec.md 中的按钮与缺失 PnL 补取任务,并完成后端测试、前端构建和浏览器验证。 + +## Comments + +- 已采用服务器固定缺失集合与已有同步任务恢复机制。 +- 实现完成:新增 pnl_backfill 任务、已提交页签按钮和任务面板提示。 +- 验证通过:25 项后端测试、Ruff、TypeScript/Vite 构建;隔离浏览器筛选为 1 条时补取 207 条,重复点击得到总数 0 的完成任务;390px 页面无横向溢出。 diff --git a/.scratch/submitted-pnl/spec.md b/.scratch/submitted-pnl/spec.md new file mode 100644 index 0000000..cf1ca3b --- /dev/null +++ b/.scratch/submitted-pnl/spec.md @@ -0,0 +1,7 @@ +# 已提交 Alpha 补取 PnL + +在已提交页签增加“检查pnl”按钮。覆盖全部本地已同步的已提交 Alpha(含隐藏记录及各地区),不受筛选、分页和勾选影响;仅补取 pnl_cache 中不存在的记录。 + +复用同步任务机制:服务端固定缺失 ID 集合,不受按 ID 操作的 100 条限制;活动任务去重,逐项落库,执行时再次检查缓存。任务支持进度、取消、连接恢复和失败重试。已有缓存不刷新,空集合直接完成并提示已齐全。无需数据库迁移。 + +验证:后端覆盖范围、超过 100 条、重复点击、已有缓存保护、逐项失败与重试、空集合和参数拒绝;浏览器验证按钮仅出现在已提交页签、筛选不影响范围、任务进度与空集合提示。 diff --git a/backend/app/business.py b/backend/app/business.py index 4a7aa8d..220e4a7 100644 --- a/backend/app/business.py +++ b/backend/app/business.py @@ -10,7 +10,7 @@ from uuid import uuid4 from fastapi import HTTPException from sqlalchemy import delete, func, select, update -from .alphas import list_statement, sorted_statement, summary +from .alphas import list_statement, sorted_statement, submission_condition, summary from .jobs import ACTIVE from .models import Account, Alpha, Job, JobItem, Pnl, Research, ResearchTag, SelfCorrelation, now from .research.provenance import alpha_sources, source_kinds @@ -203,9 +203,22 @@ class Business: for job in ( await self.db.scalars(select(Job).where(Job.kind == body.kind, Job.status.in_(ACTIVE))) ).all(): - if job.payload == payload: + if body.kind == "pnl_backfill" or job.payload == payload: return JobOutput.model_validate(job).model_dump(mode="json") job = Job(id=str(uuid4()), kind=body.kind, payload=payload) + if body.kind == "pnl_backfill": + # Fix the full missing set on the server, independently of UI paging. + # The account lock above also serializes duplicate button clicks. + ids = list((await self.db.scalars( + select(Alpha.id) + .outerjoin(Pnl, Pnl.alpha_id == Alpha.id) + .where(submission_condition("SUBMITTED"), Pnl.alpha_id.is_(None)) + .order_by(Alpha.id) + )).all()) + job.payload = {"alpha_ids": ids, "submission": "SUBMITTED"} + job.total = len(ids) + if not ids: + job.status = "completed" self.db.add(job) await self.db.flush() return JobOutput.model_validate(job).model_dump(mode="json") diff --git a/backend/app/jobs.py b/backend/app/jobs.py index 684f8b2..ec43660 100644 --- a/backend/app/jobs.py +++ b/backend/app/jobs.py @@ -412,21 +412,28 @@ class Runner: await db.commit() async def sync_ids(self, job_id, kind, alpha_ids): + """Process fixed IDs, preserving existing PnL during a missing-only backfill. + + Successful items commit individually so cancellation and retries keep + completed work. Cache presence is rechecked when a queued item runs. + """ + pnl_job = kind in ("pnl_refresh", "pnl_backfill") await self.checkpoint(job_id, {"total": len(alpha_ids)}) for alpha_id in alpha_ids: async with self.sessions() as db: previous = await db.get(JobItem, (job_id, alpha_id)) if previous and not previous.error: continue - await self.checkpoint(job_id, {"next_retry_at": None}) + cached = await db.get(Pnl, alpha_id) if kind == "pnl_backfill" else None + await self.checkpoint(job_id, {"next_retry_at": None, "checkpoint": {"alpha_id": alpha_id}}) error = None try: - raw = await ( - self.client.pnl(alpha_id) if kind == "pnl_refresh" else self.client.alpha(alpha_id) + raw = cached.raw if cached is not None else await ( + self.client.pnl(alpha_id) if pnl_job else self.client.alpha(alpha_id) ) - if kind != "pnl_refresh" and raw.get("id") != alpha_id: + if not pnl_job and raw.get("id") != alpha_id: raise ValueError("平台返回的 Alpha ID 与请求不一致") - points = pnl_points(raw) if kind == "pnl_refresh" else None + points = cached.points if cached is not None else pnl_points(raw) if pnl_job else None except VerificationRequired: raise except WqError as exc: @@ -446,10 +453,10 @@ class Runner: previous = JobItem(job_id=job_id, alpha_id=alpha_id) db.add(previous) if not error: - if kind == "pnl_refresh": + if pnl_job: if not await db.get(Alpha, alpha_id): error = "请先导入此 Alpha" - else: + elif kind != "pnl_backfill" or await db.get(Pnl, alpha_id) is None: await self.save_pnl(db, alpha_id, raw, points) else: await db.scalar(select(Account).where(Account.id == 1).with_for_update()) diff --git a/backend/app/schemas.py b/backend/app/schemas.py index 22d98a6..e53cabe 100644 --- a/backend/app/schemas.py +++ b/backend/app/schemas.py @@ -187,7 +187,7 @@ class BulkUpdate(BulkInput): class JobInput(Contract): - kind: Literal["full_sync", "daily_sync", "alpha_refresh", "pnl_refresh", "self_correlation"] + kind: Literal["full_sync", "daily_sync", "alpha_refresh", "pnl_refresh", "pnl_backfill", "self_correlation"] alpha_ids: list[str] = Field(default_factory=list) submission: Submission | None = None date_from: date | None = None @@ -195,7 +195,10 @@ class JobInput(Contract): @model_validator(mode="after") def validate_ids(self): - if self.kind in ("full_sync", "daily_sync"): + if self.kind == "pnl_backfill": + if self.alpha_ids or self.submission is not None or self.date_from is not None or self.date_to is not None: + raise ValueError("检查 PnL 自动覆盖全部本地已提交 Alpha,不接受 ID、分组或日期范围") + elif self.kind in ("full_sync", "daily_sync"): if self.alpha_ids: raise ValueError("列表同步不接受 Alpha ID") if self.kind == "full_sync": diff --git a/backend/tests/test_pnl_backfill.py b/backend/tests/test_pnl_backfill.py new file mode 100644 index 0000000..9628378 --- /dev/null +++ b/backend/tests/test_pnl_backfill.py @@ -0,0 +1,103 @@ +"""Submitted PnL backfill covers the local library and preserves existing caches.""" + +import pytest +from sqlalchemy import func, select + +from app.alphas import upsert_alpha +from app.models import Pnl +from app.worldquant import WqError +from tests.conftest import alpha +from tests.test_jobs import ready_runner, result + +URL = "/api/v1/sync-jobs" + + +async def start(client): + response = await client.post(URL, json={"kind": "pnl_backfill"}) + assert response.status_code == 202 + return response.json() + + +async def test_backfill_all_submitted_beyond_page_limit_and_deduplicates(app, logged_in): + runner = await ready_runner(app) + async with runner.sessions() as db: + for i in range(105): + await upsert_alpha(db, alpha( + f"ref{i:03}", status="DECOMMISSIONED" if i % 2 else "ACTIVE", + hidden=bool(i % 2), settings={"region": "EUR" if i % 2 else "USA"}, + )) + for raw in [alpha("pending"), alpha("unknown", status=None), alpha("cached", status="ACTIVE")]: + await upsert_alpha(db, raw) + db.add(Pnl(alpha_id="cached", raw={"original": True}, points=[])) + await db.commit() + first = await start(logged_in) + assert first["total"] == 105 + assert first["payload"]["alpha_ids"] == [f"ref{i:03}" for i in range(105)] + # Another task can populate an item after this snapshot was fixed. + async with runner.sessions() as db: + db.add(Pnl(alpha_id="ref000", raw={"original": True}, points=[])) + await db.commit() + assert (await start(logged_in))["id"] == first["id"] + calls = [] + original_pnl = runner.client.pnl + + async def pnl(alpha_id): + calls.append(alpha_id) + return await original_pnl(alpha_id) + + runner.client.pnl = pnl + await runner.execute(first["id"]) + finished = await result(runner, first["id"]) + assert (finished.status, finished.processed, finished.failed) == ("completed", 105, 0) + assert calls == [f"ref{i:03}" for i in range(1, 105)] + async with runner.sessions() as db: + assert await db.scalar(select(func.count()).select_from(Pnl)) == 106 + assert (await db.get(Pnl, "cached")).raw == {"original": True} + assert (await db.get(Pnl, "ref000")).raw == {"original": True} + empty = await start(logged_in) + assert empty["status"] == "completed" and empty["total"] == 0 + assert empty["payload"]["alpha_ids"] == [] + + +async def test_backfill_keeps_progress_and_retries_only_unfinished_items(app, logged_in): + runner = await ready_runner(app) + async with runner.sessions() as db: + for name in ("a", "b", "c"): + await upsert_alpha(db, alpha(name, status="ACTIVE")) + await db.commit() + calls = [] + blocked = True + original_pnl = runner.client.pnl + + async def pnl(alpha_id): + calls.append(alpha_id) + if blocked and alpha_id == "b": + raise WqError("无权访问", "access_denied") + if blocked and alpha_id == "c": + raise WqError("平台数据仍在准备,请稍后重试", "pending") + return await original_pnl(alpha_id) + + runner.client.pnl = pnl + job = await start(logged_in) + await runner.execute(job["id"]) + failed = await result(runner, job["id"]) + assert (failed.status, failed.processed, failed.failed) == ("failed", 1, 1) + assert failed.checkpoint["alpha_id"] == "c" + async with runner.sessions() as db: + assert await db.get(Pnl, "a") is not None + assert await db.get(Pnl, "b") is None + assert await db.get(Pnl, "c") is None + blocked = False + assert (await logged_in.post(f"{URL}/{job['id']}/retry")).status_code == 200 + await runner.execute(job["id"]) + finished = await result(runner, job["id"]) + assert (finished.status, finished.processed, finished.failed) == ("completed", 3, 0) + assert calls == ["a", "b", "c", "b", "c"] + + +@pytest.mark.parametrize("extra", [ + {"alpha_ids": ["a"]}, {"submission": "UNSUBMITTED"}, {"date_from": "2025-01-01"}, +]) +async def test_backfill_rejects_client_scope(logged_in, extra): + response = await logged_in.post(URL, json={"kind": "pnl_backfill", **extra}) + assert response.status_code == 422 diff --git a/frontend/src/api.ts b/frontend/src/api.ts index dd9d436..1ac4bdc 100644 --- a/frontend/src/api.ts +++ b/frontend/src/api.ts @@ -89,6 +89,7 @@ export const jobLabels: Record = { self_correlation: "本地自相关检测", alpha_refresh: "导入 / 刷新 Alpha", pnl_refresh: "获取 PnL", + pnl_backfill: "检查已提交 Alpha 的 PnL", connect: "连接 WorldQuant", verify: "继续人工验证", profile: "刷新个人资料", diff --git a/frontend/src/components/JobPanel.tsx b/frontend/src/components/JobPanel.tsx index 74361fc..95c5010 100644 --- a/frontend/src/components/JobPanel.tsx +++ b/frontend/src/components/JobPanel.tsx @@ -93,6 +93,13 @@ export function JobPanel({ {job.checkpoint.dates_completed} / {job.checkpoint.dates_total} 天

)} + {job.kind === "pnl_backfill" && ( +

+ {job.total === 0 + ? "已提交 Alpha 的 PnL 已齐全,无需补取" + : `仅补取缺失的 PnL${job.checkpoint?.alpha_id ? ` · ${job.checkpoint.alpha_id}` : ""}`} +

+ )} {job.kind === "self_correlation" && job.checkpoint?.alpha_id && (

{job.checkpoint.alpha_id} ·{" "} diff --git a/frontend/src/pages/AlphaPage.tsx b/frontend/src/pages/AlphaPage.tsx index 3a48618..7bf5f1d 100644 --- a/frontend/src/pages/AlphaPage.tsx +++ b/frontend/src/pages/AlphaPage.tsx @@ -33,6 +33,7 @@ import type { Alpha, AlphaPage as Page, Facets, + Job, Submission, } from "../types"; import { AlphaDetail } from "../components/AlphaDetail"; @@ -313,13 +314,15 @@ export function AlphaPage({ async function newTask(kind: string, ids: string[] = []) { setBusy(kind); try { - await post("/sync-jobs", { + const job = await post("/sync-jobs", { kind, alpha_ids: ids, ...(kind === "full_sync" ? { submission: "SUBMITTED" } : {}), }); setImporting(false); setIdText(""); + if (kind === "pnl_backfill" && job.total === 0) + Toast.success("已提交 Alpha 的 PnL 已齐全,无需补取"); onTask(); } catch (e) { Toast.error((e as Error).message); @@ -966,6 +969,17 @@ export function AlphaPage({ > 按天同步 + {submission === "SUBMITTED" && ( + + )} {submission === "SUBMITTED" && (