feat: backfill missing PnL for submitted alphas
This commit is contained in:
+15
-2
@@ -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")
|
||||
|
||||
+14
-7
@@ -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())
|
||||
|
||||
@@ -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":
|
||||
|
||||
Reference in New Issue
Block a user