diff --git a/.scratch/alpha-management/issues/01-implement.md b/.scratch/alpha-management/issues/01-implement.md new file mode 100644 index 0000000..31c0445 --- /dev/null +++ b/.scratch/alpha-management/issues/01-implement.md @@ -0,0 +1,18 @@ +# 实现 Alpha 分组同步与本地自相关 + +Status: ready-for-agent + +## 工作项 + +- [x] 分组查询、按日/已提交全量同步契约与持久化进度。 +- [x] 本地自相关算法、缓存补取、结果存储与数据不足处理。 +- [x] 双 Tab、同步日期表单、检测入口与详情、任务进度。 +- [x] 后端、迁移、前端构建及浏览器验证,更新当前项目使用文档。 + +## Comments + +2026-09-08:开始实现;用户已明确授权上述范围内的本地修改。 + +2026-09-08:本地实现完成。`uv run ruff check app tests` 与后端 87 项测试通过;`pnpm build`、修改文件 Prettier 检查和浏览器全套 6 项通过。截图核验后修正检测结果与缓存时间的 UTC 标记,后端全套及 Alpha 管理浏览器流程再次通过。已检查 1440px / 390px 布局。构建仅有已有 lottie-web 依赖的 eval 提示。 + +在隔离 SQLite 与 PostgreSQL 17 中验证 0002 → 0003、Alembic schema check、回退再升级,旧研究记录内容保留;临时 PostgreSQL 容器及卷已清理。独立后端核验未发现实质问题。`git diff --check` 通过。未提交、未部署、未调用真实平台;新增平台日期筛选参数仍需真实账户只读联调。 diff --git a/.scratch/alpha-management/spec.md b/.scratch/alpha-management/spec.md new file mode 100644 index 0000000..2cabd77 --- /dev/null +++ b/.scratch/alpha-management/spec.md @@ -0,0 +1,25 @@ +# Alpha 管理迭代 + +用户于 2026-09-08 要求实现本地自相关检测、待提交/已提交双 Tab,以及按日同步。 + +## 范围与行为 + +- Alpha 仍存于统一结果库,按平台 `status == UNSUBMITTED` 与已知的其他状态分成两个 Tab;筛选、导出、选择及 AI 页面上下文携带所属分组。 +- 待提交先选日期范围,按 UTC 创建日期逐天同步;已提交按 UTC 提交日期逐天同步,另可全量同步。相同起止日期即单日。覆盖可见及隐藏数据,半开日边界防止午夜遗漏;保留分页检查点、取消、重试与本地研究记录。 +- 新建全量任务只允许已提交。升级前已有的无分组全量任务按原始任务范围恢复,不改写其检查点含义。 +- 本地 self-correlation:同地区已提交 Alpha 作基准,排除自身;使用 PnL 缓存并在缺失时按需补取,只计算本地 Pearson,不调用平台 check。 +- 累计 PnL 先按 UTC 日期排序并作日变化,缺失值不补零、不跨缺口差分;以目标最新日期为四年共同回看窗口;至少 30 个共同有效变化样本,常量或无效数据跳过。 +- 保留带符号最大相关系数和 0.7 告警线(本地规则,不宣称平台等价)。展示比较数量、跳过原因、最高相关对象、计算时间和缓存时间;缺失数据不得显示通过。 +- 检测结果单独持久化,不覆盖平台 checks 或本地研究状态。缓存和基准集改变时标为待重算。单条详情和选中最多 100 条均可发起任务。 + +## 界面约定 + +scope_sketch: 现有 Alpha 管理中增加两个 Tab、同步日期对话框及本地自相关详情;不增加其他研究模块。 +lark_style_recipe: 保留白色工作区、紧凑表格、4px 间距和克制蓝色主操作;不改变页面外壳。 +ud_control_coverage: 复用 Semi Design Tabs、Input、Modal、Button、Table、Tag 的语义与状态。 +media_decision: 数据管理任务不需要新增插图或图标。 +verification: 单元/接口验证数值、日期边界、任务恢复和状态隔离;浏览器验证双 Tab、同步范围、检测、导出、窄屏和现有 AI 流程。 + +## 验证边界 + +只在隔离测试数据库和模拟平台上验证;不调用真实平台检查,不启动回测,不回写平台。 diff --git a/README.md b/README.md index 32c0319..5ec5c85 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # WorldQuant Alpha 研究工作空间 -个人单账户系统。首期实现平台资料、Alpha 同步与查询、PnL 缓存、本地备注/标签/收藏/研究状态。平台接口只读,认证除外;不会回测、触发检查、回写属性或提交 Alpha。 +个人单账户系统。实现平台资料、Alpha 分组同步与查询、PnL 缓存、本地自相关检测、本地备注/标签/收藏/研究状态。平台接口只读,认证除外;不会回测、触发平台检查、回写属性或提交 Alpha。 需求与后续路线图见 [项目方案](docs/project-plan.md),AI 助手范围见 [开发计划](docs/ai-chatbot-plan.md)。前端 React 19 + TypeScript + Semi Design,后端 Python 3.12 + FastAPI + HTTPX + SQLAlchemy,PostgreSQL 保存数据,Caddy 提供 Web 入口。前后端独立依赖、独立构建,所有部署文件位于根目录。 @@ -22,8 +22,10 @@ docker compose ps 1. 登录系统,在“个人信息”保存 WorldQuant 邮箱和密码,点击“连接 WorldQuant”。 2. 如平台要求人工验证,在显示的入口完成操作,再点击“继续验证”;后台保留同一挑战会话。 -3. 在“Alpha 管理”手动同步平台数据,或导入指定 Alpha ID。任务面板显示进度、错误、取消和重试。 -4. 点击 Alpha 打开详情。研究记录保存在本地;PnL 点击获取后缓存。下次同步会更新平台数据并保留本地研究记录。 +3. 在“Alpha 管理”切换“待提交 / 已提交”。待提交先选创建日期范围,按天同步;已提交可选提交日期范围按天同步,或全量同步。相同起止日期表示单日,日期边界为 UTC,均包含隐藏记录。也可导入指定 Alpha ID。任务面板显示当前日期、进度、错误、取消和重试。 +4. 点击 Alpha 打开详情。研究记录保存在本地;PnL 点击获取后缓存。详情的“本地自相关”可发起检测,列表也可选中最多 100 条批量检测。建议先全量同步已提交 Alpha,建立比较基准。下次同步会更新平台数据并保留本地研究记录。 + +本地自相关只与本地已同步、同地区的已提交 Alpha 比较,排除自身。使用累计 PnL 的日变化,在目标最新数据日往前四年的共同窗口计算 Pearson 相关系数,至少需要 30 个共同有效样本;带符号最大值达到 0.7 时提示相关性偏高。这是本地规则,不等同于平台检查。PnL 缓存缺失时自动补取;样本不足、常量序列和不可用基准会明确显示,不能当作通过。结果独立保存,相关 PnL、地区或基准成员变化后标记为“待重算”。 平台未返回的资料与指标保留为空。数值筛选采用平台原始单位,例如 Turnover `0.15` 表示 15%。日期筛选边界为 UTC;时间显示采用个人页的时区偏好。 @@ -100,7 +102,7 @@ docker compose logs --tail=100 backend ## 备份与恢复 -以下为本机配置命令;公网统一补上 `-f compose.public.yaml`,自定义项目名时保持相同 `-p`。数据库备份包括平台快照、研究记录及版本、账户密文、模型配置密文、AI 会话/消息/执行/工具确认记录和同步任务。备份文件仍属于私有数据。 +以下为本机配置命令;公网统一补上 `-f compose.public.yaml`,自定义项目名时保持相同 `-p`。数据库备份包括平台快照、研究记录及版本、本地自相关结果、账户密文、模型配置密文、AI 会话/消息/执行/工具确认记录和同步任务。备份文件仍属于私有数据。 ```bash mkdir -p backups @@ -148,7 +150,7 @@ pnpm exec playwright install chromium pnpm test ``` -浏览器测试自动启动临时数据库、模拟平台 API 和 Vite,使用 620 条明确标记 `TEST` 的合成 Alpha。不会向正式数据库写入样例。测试验证模型配置、查询卡片、修改预览及确认、草稿冲突、收起及刷新恢复、取消,以及系统登录、账户连接、多页同步、SUPER 详情、备注与收藏在刷新后保留、PnL、超过 500 条 CSV 及退出。截图写入忽略目录 `output/playwright/`。 +浏览器测试自动启动临时数据库、模拟平台 API 和 Vite,使用 620 条明确标记 `TEST` 的合成 Alpha。不会向正式数据库写入样例。测试验证模型配置、查询卡片、修改预览及确认、草稿冲突、收起及刷新恢复、取消,以及系统登录、账户连接、双 Tab、按日及全量同步、本地自相关保存、SUPER 详情、备注与收藏保留、PnL、分组 CSV 导出及退出。截图写入忽略目录 `output/playwright/`。 需要重跑 Docker 持久化和备份验收时,创建独立测试环境,并在测试 env 文件中选择空闲 `LOCAL_PORT`(例如 18089),无需停止正式实例。以下脚本只接受 `wq-alpha-acceptance*` 项目名: @@ -182,14 +184,15 @@ FastAPI 的 `/openapi.json` 与 `/docs` 可在后端开发端口访问;生产 - `/api/v1/account`:偏好、加密凭据、连接/验证/断开/资料刷新。 - `/api/v1/alphas`:服务端筛选与排序、详情、本地研究记录、批量编辑、流式 CSV。 - `/api/v1/alphas/{id}/pnl`:只读缓存;刷新通过 `pnl_refresh` 任务。 +- `/api/v1/alphas/{id}/self-correlation`:读取本地检测结果;检测通过 `self_correlation` 任务。 - `/api/v1/sync-jobs`:创建任务立即返回 202 和 ID,查询、取消与重试。 - `/api/v1/ai`:脱敏模型配置与测试、会话历史、SSE 执行、执行快照、取消及确认。新执行只接收 `request_id`、`message`、`context`;同一会话重复请求 ID 返回原运行,参数变化返回 409。 研究记录 PATCH 现在必须提供读取时的 `version`;批量编辑必须提供每个目标 ID 的 `versions` 映射。`0002` 迁移给旧研究记录设置初始版本 1,不修改其内容。版本冲突返回 409。 -写请求需 `X-WQ-Request: 1`;浏览器跨站写入被拒绝。Alpha 平台快照、`research` 本地研究、`pnl_cache` 分开存储。研究状态固定为 `inbox/candidate/optimizing/archived`;平台类型、语言、状态按原值显示。 +写请求需 `X-WQ-Request: 1`;浏览器跨站写入被拒绝。Alpha 平台快照、`research` 本地研究、`pnl_cache`、`self_correlations` 本地检测结果分开存储。`0003` 迁移只新增检测结果表。研究状态固定为 `inbox/candidate/optimizing/archived`;平台类型、语言、状态按原值显示。 -同步按“未提交/已提交 × 可见/隐藏”分页,每页数据与检查点同事务提交,Alpha ID 幂等更新。失败任务保留进度,重试只处理剩余页或失败 ID。上游 `Retry-After` 等待可被取消。分页过程中平台记录移动可能造成重复或遗漏,通过 ID 去重和再次全量同步校正;单次没有查到不自动删除本地记录。 +列表及导出支持 `submission=UNSUBMITTED|SUBMITTED`,平台状态缺失时不推断为已提交。`daily_sync` 必须提供分组及 `date_from` / `date_to`,每个 UTC 日期分别分页获取可见、隐藏记录;新建 `full_sync` 只同步已提交。旧的无分组全量任务保持原范围恢复。每页数据与检查点同事务提交,Alpha ID 幂等更新。失败任务保留进度,重试只处理剩余页或失败 ID。上游 `Retry-After` 等待可被取消。分页过程中平台记录移动可能造成重复或遗漏,通过 ID 去重和再次同步对应范围校正;单次没有查到不自动删除本地记录。 ## 日志排查与验证边界 @@ -205,4 +208,4 @@ curl -f http://localhost:8080/api/v1/health AI 模型兼容性由模拟 Chat Completions/Responses HTTP 流与真实 SDK 适配器验证;未配置真实供应商前,不能保证其工具选择质量、模型权限或网关兼容性。真实联调请分别记录流式回答与业务工具调用是否成功。 -实现使用旧项目已知请求形态并对模拟上游做自动化验证。WorldQuant 当前真实账号权限、人工验证页面行为、实际数据 schema、真实账户全量同步及公网证书签发,均需要在自己的账户/域名完成只读联调;未取得该证据前不宣称已验证。验收实测结果见 [验收记录](docs/verification.md)。 +实现参考旧项目请求形态,并对模拟上游做自动化验证。新增日期筛选参数、WorldQuant 当前真实账号权限、人工验证页面行为、实际数据 schema、真实账户同步及公网证书签发,均需要在自己的账户/域名完成只读联调;未取得该证据前不宣称已验证。验收实测结果见 [验收记录](docs/verification.md)。 diff --git a/backend/app/ai/tools.py b/backend/app/ai/tools.py index 6f5446e..2c18d84 100644 --- a/backend/app/ai/tools.py +++ b/backend/app/ai/tools.py @@ -60,7 +60,7 @@ CATALOG = { "bulk_update_research": (BulkInput, "提出固定 1–100 个 Alpha 的批量标签或研究状态修改,等待用户确认。"), "create_sync_job": ( JobInput, - "提出全量同步、指定 Alpha 刷新或 PnL 刷新任务,等待确认;创建后立即返回任务 ID。", + "提出同步或本地自相关任务,等待确认。full_sync 仅同步已提交;待提交必须用 daily_sync 并指定 submission、date_from/date_to(UTC),待提交按创建日、已提交按提交日逐天同步。alpha_refresh/pnl_refresh/self_correlation 使用固定 alpha_ids;自相关缺失 PnL 时自动补取,不触发平台检查。创建后立即返回任务 ID。", ), "cancel_job": (JobArgs, "提出取消指定同步任务,等待用户确认。"), "retry_job": (JobArgs, "提出重试指定失败或暂停的同步任务,等待用户确认。"), diff --git a/backend/app/alphas.py b/backend/app/alphas.py index 2f96604..9ddf609 100644 --- a/backend/app/alphas.py +++ b/backend/app/alphas.py @@ -4,9 +4,24 @@ import math import re from datetime import datetime -from sqlalchemy import or_, select +from sqlalchemy import or_, select, update + +from .models import Alpha, Research, ResearchTag, SelfCorrelation, now + + +def submission_condition(submission): + """Match the platform list contract; a missing status is never assumed submitted.""" + return Alpha.status == "UNSUBMITTED" if submission == "UNSUBMITTED" else Alpha.status != "UNSUBMITTED" + + +async def invalidate_correlations(db, alpha_id, regions=()): + """A changed baseline or PnL invalidates local conclusions without touching platform checks.""" + await db.execute( + update(SelfCorrelation) + .where(or_(SelfCorrelation.alpha_id == alpha_id, SelfCorrelation.region.in_(regions))) + .values(stale=True) + ) -from .models import Alpha, Research, ResearchTag, now SENSITIVE_KEYS = { "password", @@ -67,6 +82,8 @@ async def upsert_alpha(db, raw: dict): if not isinstance(alpha_id, str) or not alpha_id: raise ValueError("Alpha 数据缺少 ID") item = await db.get(Alpha, alpha_id) + previous_region = item.region if item else None + previous_status = item.status if item else None if item is None: item = Alpha(id=alpha_id) db.add(item) @@ -78,6 +95,13 @@ async def upsert_alpha(db, raw: dict): item.alpha_type, item.language = raw.get("type"), settings.get("language") item.stage, item.status, item.hidden = raw.get("stage"), raw.get("status"), raw.get("hidden") is True item.region, item.universe = settings.get("region"), settings.get("universe") + if previous_region != item.region or previous_status != item.status: + regions = { + region + for region, status in ((previous_region, previous_status), (item.region, item.status)) + if region and status and status != "UNSUBMITTED" + } + await invalidate_correlations(db, alpha_id, regions) item.settings, item.is_metrics = sanitize(settings), sanitize(metrics) item.os_metrics = sanitize(raw.get("os")) if isinstance(raw.get("os"), dict) else {} item.checks = sanitize(metrics.get("checks") or raw.get("checks") or []) @@ -93,6 +117,8 @@ async def upsert_alpha(db, raw: dict): def list_statement(filters): query = select(Alpha, Research).join(Research, Research.alpha_id == Alpha.id) + if filters.submission: + query = query.where(submission_condition(filters.submission)) q = filters.q if q: pattern = "%" + q.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%" diff --git a/backend/app/business.py b/backend/app/business.py index c63c362..db6dba9 100644 --- a/backend/app/business.py +++ b/backend/app/business.py @@ -4,6 +4,7 @@ Mutations never commit here, so the AI executor can atomically save their audit Job runner notifications must happen after commit, using ``notify_job``. """ +from datetime import timezone from uuid import uuid4 from fastapi import HTTPException @@ -11,7 +12,7 @@ from sqlalchemy import delete, func, select, update from .alphas import list_statement, sorted_statement, summary from .jobs import ACTIVE -from .models import Account, Alpha, Job, JobItem, Pnl, Research, ResearchTag, now +from .models import Account, Alpha, Job, JobItem, Pnl, Research, ResearchTag, SelfCorrelation, now from .schemas import AlphaDetail, AlphaPage, BulkUpdate, JobInput, JobOutput, ResearchUpdate, normalize_tags @@ -29,10 +30,51 @@ class Business: .offset(filters.offset) ) ).all() + correlations = { + row.alpha_id: self.correlation_summary(row) + for row in ( + await self.db.scalars( + select(SelfCorrelation).where(SelfCorrelation.alpha_id.in_([a.id for a, _ in rows])) + ) + ).all() + } return AlphaPage( - items=[summary(a, r) for a, r in rows], total=total, limit=filters.limit, offset=filters.offset + items=[{**summary(a, r), "local_correlation": correlations.get(a.id)} for a, r in rows], + total=total, + limit=filters.limit, + offset=filters.offset, ).model_dump(mode="json") + @staticmethod + def correlation_summary(row): + return { + **{ + key: row.result.get(key) + for key in ("status", "max_correlation", "compared_count", "skipped_count") + }, + "stale": row.stale, + "calculated_at": row.calculated_at.replace( + tzinfo=row.calculated_at.tzinfo or timezone.utc + ).isoformat(), + } + + async def get_self_correlation(self, alpha_id): + if not await self.db.get(Alpha, alpha_id): + raise HTTPException(404, "Alpha 尚未同步") + row = await self.db.get(SelfCorrelation, alpha_id) + return { + "cached": row is not None, + "result": { + **row.result, + "stale": row.stale, + "calculated_at": row.calculated_at.replace( + tzinfo=row.calculated_at.tzinfo or timezone.utc + ).isoformat(), + } + if row + else None, + } + async def get_alpha_facets(self): result = {} for key in ("region", "universe", "alpha_type", "language", "status", "stage"): @@ -126,9 +168,15 @@ class Business: async def create_sync_job(self, body: JobInput): account = await self.db.scalar(select(Account).where(Account.id == 1).with_for_update()) - if not account.password_encrypted or account.connection_status in ("disconnected", "error"): + if body.kind != "self_correlation" and ( + not account.password_encrypted or account.connection_status in ("disconnected", "error") + ): raise HTTPException(409, "请先连接 WorldQuant") - payload = {"alpha_ids": body.alpha_ids} + if body.kind == "self_correlation": + found = set((await self.db.scalars(select(Alpha.id).where(Alpha.id.in_(body.alpha_ids)))).all()) + if found != set(body.alpha_ids): + raise HTTPException(404, "部分 Alpha 尚未同步") + payload = body.model_dump(mode="json", exclude={"kind"}, exclude_none=True) for job in ( await self.db.scalars(select(Job).where(Job.kind == body.kind, Job.status.in_(ACTIVE))) ).all(): diff --git a/backend/app/correlation.py b/backend/app/correlation.py new file mode 100644 index 0000000..0e6e77f --- /dev/null +++ b/backend/app/correlation.py @@ -0,0 +1,138 @@ +"""Local Pearson comparison of daily PnL changes; no platform eligibility decisions. + +The four-year window and signed 0.7 threshold follow the legacy research tool. +Thirty paired changes is a local minimum, not a claim about BRAIN's own checks. +""" + +import math +from datetime import datetime, timezone +from statistics import StatisticsError, correlation + +THRESHOLD = 0.7 +MIN_SAMPLES = 30 +WINDOW_YEARS = 4 + + +def daily_changes(points): + """Return dated changes and the final date, rejecting ambiguous daily records. + + Parameters are normalized PnL points. Missing values break a change interval; + each change retains its start date so differently spaced samples never pair. + Raises ValueError for malformed dates, duplicate days or non-finite values. + """ + daily = {} + for point in points: + try: + timestamp = datetime.fromisoformat(point["date"].replace("Z", "+00:00")) + day = ( + (timestamp.replace(tzinfo=timezone.utc) if timestamp.tzinfo is None else timestamp) + .astimezone(timezone.utc) + .date() + ) + except (ValueError, TypeError, KeyError, AttributeError): + raise ValueError("PnL 日期无法识别") from None + if day in daily: + raise ValueError("PnL 同一天存在多条记录") + value = point.get("value") + if value is not None and ( + isinstance(value, bool) or not isinstance(value, (int, float)) or not math.isfinite(value) + ): + raise ValueError("PnL 包含无效数值") + daily[day] = value + changes, previous_day, previous = {}, None, None + for day, value in sorted(daily.items()): + if value is not None and previous is not None: + delta = value - previous + if not math.isfinite(delta): + raise ValueError("PnL 变化超出有效数值范围") + changes[day] = (previous_day, delta) + previous_day, previous = day, value + return changes, max(daily) if daily else None + + +def calculate_correlation(target_points, references): + """Compare a target with supplied reference caches and return a JSON-safe report. + + Each reference has alpha_id, points, fetched_at and optionally error. Callers + select the same-region submitted set and exclude the target. Incomplete + coverage can report high correlation, but can never report an all-clear. + """ + result = { + "status": "insufficient_data", + "threshold": THRESHOLD, + "min_samples": MIN_SAMPLES, + "window_years": WINDOW_YEARS, + "window_from": None, + "window_to": None, + "max_correlation": None, + "most_correlated_alpha_id": None, + "candidate_count": len(references), + "compared_count": 0, + "skipped_count": 0, + "matches": [], + "skipped": [], + "reason": None, + } + try: + target, latest = daily_changes(target_points) + except ValueError as exc: + result["reason"] = str(exc) + return result + if latest is None or not target: + result["reason"] = "目标 Alpha 没有可用的 PnL 日变化" + return result + try: + cutoff = latest.replace(year=latest.year - WINDOW_YEARS) + except ValueError: + cutoff = latest.replace(year=latest.year - WINDOW_YEARS, day=28) + result["window_from"], result["window_to"] = cutoff.isoformat(), latest.isoformat() + matches, skipped = [], [] + for reference in references: + alpha_id, reason = reference["alpha_id"], reference.get("error") + if not reason: + try: + changes, _ = daily_changes(reference["points"]) + days = sorted( + day + for day in target.keys() & changes.keys() + if cutoff < day <= latest and target[day][0] == changes[day][0] + ) + if len(days) < MIN_SAMPLES: + reason = f"共同有效样本不足 {MIN_SAMPLES} 个(实际 {len(days)})" + else: + coefficient = correlation([target[d][1] for d in days], [changes[d][1] for d in days]) + if not math.isfinite(coefficient): + reason = "无法计算有效相关系数" + else: + matches.append( + { + "alpha_id": alpha_id, + "correlation": coefficient, + "sample_count": len(days), + "date_from": days[0].isoformat(), + "date_to": days[-1].isoformat(), + "pnl_fetched_at": reference.get("fetched_at"), + } + ) + except StatisticsError: + reason = "目标或基准 PnL 日变化为常量" + except (ValueError, OverflowError) as exc: + reason = str(exc) + if reason: + skipped.append({"alpha_id": alpha_id, "reason": reason}) + matches.sort(key=lambda row: (-row["correlation"], row["alpha_id"])) + result.update( + compared_count=len(matches), skipped_count=len(skipped), matches=matches[:10], skipped=skipped[:100] + ) + if matches: + maximum = matches[0]["correlation"] + result.update( + max_correlation=maximum, + most_correlated_alpha_id=matches[0]["alpha_id"], + status="high" if maximum >= THRESHOLD else "partial" if skipped else "low", + ) + else: + result["reason"] = ( + "没有同地区已提交 Alpha 可供比较" if not references else "所有基准均缺少足够的有效样本" + ) + return result diff --git a/backend/app/jobs.py b/backend/app/jobs.py index fa369b5..4fe5af4 100644 --- a/backend/app/jobs.py +++ b/backend/app/jobs.py @@ -8,14 +8,15 @@ a distributed lease and session coordinator. import asyncio import logging from contextlib import suppress -from datetime import timedelta +from datetime import date, datetime, time, timedelta, timezone from uuid import uuid4 from sqlalchemy import select, update from sqlalchemy.exc import SQLAlchemyError -from .alphas import pnl_points, sanitize, upsert_alpha -from .models import Account, Alpha, Job, JobItem, Pnl, now +from .alphas import invalidate_correlations, pnl_points, sanitize, submission_condition, upsert_alpha +from .correlation import calculate_correlation +from .models import Account, Alpha, Job, JobItem, Pnl, SelfCorrelation, now from .security import cipher from .worldquant import VerificationRequired, WqClient, WqError @@ -235,7 +236,10 @@ class Runner: async with self.sessions() as db: job = await db.get(Job, job_id) kind, payload = job.kind, job.payload - if kind == "verify": + if kind == "self_correlation": + # Cached local comparisons also work while the platform is disconnected. + await self.check_correlations(job_id, payload["alpha_ids"]) + elif kind == "verify": if not self.client.verification_url: await self.ensure_connected(force=True) else: @@ -245,7 +249,7 @@ class Runner: await self.ensure_connected(force=kind == "connect") if kind in ("connect", "profile"): await self.refresh_profile() - elif kind == "full_sync": + elif kind in ("full_sync", "daily_sync"): await self.sync_all(job_id) else: await self.sync_ids(job_id, kind, payload["alpha_ids"]) @@ -296,25 +300,68 @@ class Runner: self.client.on_retry = None async def sync_all(self, job_id): - partitions = [ - ("UNSUBMITTED", False), - ("UNSUBMITTED", True), - ("SUBMITTED", False), - ("SUBMITTED", True), - ] async with self.sessions() as db: job = await db.get(Job, job_id) checkpoint = job.checkpoint before = job.created_at.isoformat() + payload, daily = job.payload, job.kind == "daily_sync" + days = [None] + if daily: + first, last = date.fromisoformat(payload["date_from"]), date.fromisoformat(payload["date_to"]) + days = [first + timedelta(days=i) for i in range((last - first).days + 1)] + # Preserve the scope and partition indexes of pre-upgrade queued jobs. + # New JobInput always fixes a full sync to SUBMITTED. + submissions = [payload["submission"]] if payload.get("submission") else ["UNSUBMITTED", "SUBMITTED"] + partitions = [ + (submission, hidden, day) + for day in days + for submission in submissions + for hidden in (False, True) + ] start_partition, offset = checkpoint.get("partition", 0), checkpoint.get("offset", 0) for partition in range(start_partition, len(partitions)): - submission, hidden = partitions[partition] + submission, hidden, day = partitions[partition] + day_params = {} + if day: + day_params = { + "date_from": datetime.combine(day, time.min, timezone.utc).isoformat(), + "date_to": datetime.combine(day + timedelta(days=1), time.min, timezone.utc).isoformat(), + } while True: - await self.checkpoint(job_id, {"next_retry_at": None}) - raw = await self.client.alphas(submission, hidden, offset, before) + progress = {"partition": partition, "offset": offset} + if daily: + progress.update( + date=day.isoformat(), dates_completed=partition // 2, dates_total=len(days) + ) + await self.checkpoint(job_id, {"next_retry_at": None, "checkpoint": progress}) + raw = await self.client.alphas(submission, hidden, offset, before, **day_params) rows = raw.get("results") if not isinstance(rows, list): raise WqError("Alpha 列表缺少 results,已保留当前进度", "invalid_response") + if payload.get("submission"): + for raw_alpha in rows: + status = raw_alpha.get("status") + if not status or (status == "UNSUBMITTED") != (submission == "UNSUBMITTED"): + raise WqError( + "平台返回的 Alpha 不属于请求的提交分组,已保留进度", "invalid_response" + ) + if day: + field = "dateCreated" if submission == "UNSUBMITTED" else "dateSubmitted" + try: + timestamp = datetime.fromisoformat(raw_alpha[field].replace("Z", "+00:00")) + timestamp = ( + timestamp.replace(tzinfo=timezone.utc) + if timestamp.tzinfo is None + else timestamp + ) + matches_day = timestamp.astimezone(timezone.utc).date() == day + except (ValueError, KeyError, AttributeError, TypeError): + matches_day = False + if not matches_day: + raise WqError( + "平台未按请求的日期返回 Alpha,已保留进度,请核对平台日期筛选支持", + "invalid_response", + ) async with self.sessions() as db: job = await db.get(Job, job_id) if job.cancel_requested: @@ -335,9 +382,12 @@ class Runner: raise WqError("平台分页未前进", "invalid_response") offset += len(rows) job.checkpoint = { + **progress, "partition": partition if more else partition + 1, "offset": offset if more else 0, } + if daily: + job.checkpoint["dates_completed"] = job.checkpoint["partition"] // 2 job.updated_at = now() await db.commit() if not more: @@ -389,11 +439,7 @@ class Runner: if not await db.get(Alpha, alpha_id): error = "请先导入此 Alpha" else: - pnl = await db.get(Pnl, alpha_id) - if pnl is None: - pnl = Pnl(alpha_id=alpha_id) - db.add(pnl) - pnl.raw, pnl.points, pnl.fetched_at = sanitize(raw), points, now() + await self.save_pnl(db, alpha_id, raw, points) else: await upsert_alpha(db, raw) previous.error = error @@ -403,3 +449,156 @@ class Runner: job.processed += 1 job.updated_at = now() await db.commit() + + async def save_pnl(self, db, alpha_id, raw, points): + """Persist a cache and invalidate conclusions that depend on the changed series.""" + alpha = await db.get(Alpha, alpha_id) + pnl = await db.get(Pnl, alpha_id) + if pnl is None or pnl.points != points: + regions = ( + [alpha.region] if alpha.region and alpha.status and alpha.status != "UNSUBMITTED" else [] + ) + await invalidate_correlations(db, alpha_id, regions) + if pnl is None: + pnl = Pnl(alpha_id=alpha_id) + db.add(pnl) + pnl.raw, pnl.points, pnl.fetched_at = sanitize(raw), points, now() + return pnl + + async def correlation_pnl(self, job_id, alpha_id): + """Use a local cache, filling a missing one through the read-only adapter.""" + await self.checkpoint(job_id, {"next_retry_at": None}) + async with self.sessions() as db: + pnl = await db.get(Pnl, alpha_id) + if pnl is not None: + return { + "alpha_id": alpha_id, + "points": pnl.points, + # SQLite drops the offset; stored datetimes are still UTC. + "fetched_at": pnl.fetched_at.replace( + tzinfo=pnl.fetched_at.tzinfo or timezone.utc + ).isoformat(), + } + await self.ensure_connected() + raw = await self.client.pnl(alpha_id) + points = pnl_points(raw) + async with self.sessions() as db: + if (await db.get(Job, job_id)).cancel_requested: + raise asyncio.CancelledError() + pnl = await self.save_pnl(db, alpha_id, raw, points) + await db.commit() + return {"alpha_id": alpha_id, "points": points, "fetched_at": pnl.fetched_at.isoformat()} + + async def check_correlations(self, job_id, alpha_ids): + """Check fixed targets against same-region submitted caches, with per-target recovery. + + Missing references are reported as incomplete coverage. Authentication, + transient upstream errors and cancellation retain the task for retry. + """ + await self.checkpoint(job_id, {"total": len(alpha_ids)}) + reference_caches = {} + 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 + alpha = await db.get(Alpha, alpha_id) + region = alpha.region if alpha else None + reference_ids = ( + list( + ( + await db.scalars( + select(Alpha.id) + .where( + submission_condition("SUBMITTED"), + Alpha.region == region, + Alpha.id != alpha_id, + ) + .order_by(Alpha.id) + ) + ).all() + ) + if region + else [] + ) + await self.checkpoint( + job_id, + { + "checkpoint": { + "alpha_id": alpha_id, + "phase": "pnl", + "references_total": len(reference_ids), + "references_loaded": 0, + } + }, + ) + error, result = None, None + try: + if not alpha: + raise ValueError("Alpha 尚未同步") + if not region: + result = calculate_correlation([], []) + result["reason"] = "平台未提供目标 Alpha 的地区" + elif not reference_ids: + result = calculate_correlation([], []) + result["reason"] = "没有同地区已提交 Alpha 可供比较,请先同步已提交 Alpha" + else: + target = await self.correlation_pnl(job_id, alpha_id) + references = [] + for index, reference_id in enumerate(reference_ids): + if reference_id not in reference_caches: + try: + reference_caches[reference_id] = await self.correlation_pnl( + job_id, reference_id + ) + except WqError as exc: + if exc.code not in ("not_found", "access_denied"): + raise + reference_caches[reference_id] = {"alpha_id": reference_id, "error": str(exc)} + except ValueError as exc: + reference_caches[reference_id] = {"alpha_id": reference_id, "error": str(exc)} + references.append(reference_caches[reference_id]) + await self.checkpoint( + job_id, + { + "checkpoint": { + "alpha_id": alpha_id, + "phase": "pnl", + "references_total": len(reference_ids), + "references_loaded": index + 1, + } + }, + ) + await self.checkpoint( + job_id, {"checkpoint": {"alpha_id": alpha_id, "phase": "calculating"}} + ) + result = await asyncio.to_thread(calculate_correlation, target["points"], references) + result["target_pnl_fetched_at"] = target["fetched_at"] + except WqError as exc: + if exc.code not in ("not_found", "access_denied"): + raise + error = str(exc) + except ValueError as exc: + error = str(exc) + async with self.sessions() as db: + job = await db.get(Job, job_id) + if job.cancel_requested: + raise asyncio.CancelledError() + previous = await db.get(JobItem, (job_id, alpha_id)) + if previous and previous.error: + job.failed -= 1 + if not previous: + previous = JobItem(job_id=job_id, alpha_id=alpha_id) + db.add(previous) + previous.error = error + if error: + job.failed += 1 + else: + row = await db.get(SelfCorrelation, alpha_id) + if row is None: + row = SelfCorrelation(alpha_id=alpha_id) + db.add(row) + row.region, row.result, row.calculated_at, row.stale = region, result, now(), False + job.processed += 1 + job.updated_at = now() + await db.commit() diff --git a/backend/app/main.py b/backend/app/main.py index c71ebd2..8d92190 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -40,6 +40,7 @@ from .schemas import ( PnlOutput, PreferencesInput, ResearchUpdate, + SelfCorrelationOutput, SessionOutput, ) from .security import bootstrap, cipher, issue_session, require_auth, token_hash, valid_password @@ -341,6 +342,11 @@ def create_app(settings=None, wq_client=None, ai_model_factory=None): async with sessions() as db: return await Business(db).get_alpha_pnl(alpha_id) + @api.get("/alphas/{alpha_id}/self-correlation", response_model=SelfCorrelationOutput, tags=["alphas"]) + async def self_correlation(alpha_id: str): + async with sessions() as db: + return await Business(db).get_self_correlation(alpha_id) + @api.post("/sync-jobs", status_code=202, response_model=JobOutput, tags=["sync-jobs"]) async def new_job(body: JobInput): async with sessions.begin() as db: diff --git a/backend/app/models.py b/backend/app/models.py index 42d0868..af48f7e 100644 --- a/backend/app/models.py +++ b/backend/app/models.py @@ -111,6 +111,15 @@ class ResearchTag(Base): tag: Mapped[str] = mapped_column(String(60), primary_key=True, index=True) +class SelfCorrelation(Base): + __tablename__ = "self_correlations" + alpha_id: Mapped[str] = mapped_column(ForeignKey("alphas.id"), primary_key=True) + region: Mapped[str | None] = mapped_column(String(50), index=True) + result: Mapped[dict] = mapped_column(JSON) + stale: Mapped[bool] = mapped_column(Boolean, default=False) + calculated_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), default=now) + + class Job(Base): __tablename__ = "sync_jobs" id: Mapped[str] = mapped_column(String(36), primary_key=True) diff --git a/backend/app/schemas.py b/backend/app/schemas.py index a206ef6..db24342 100644 --- a/backend/app/schemas.py +++ b/backend/app/schemas.py @@ -1,13 +1,14 @@ """Validated public API contracts. Platform state is intentionally not a closed enum.""" import re -from datetime import datetime, timezone +from datetime import date, datetime, timezone from typing import Literal from zoneinfo import ZoneInfo, ZoneInfoNotFoundError from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator ResearchState = Literal["inbox", "candidate", "optimizing", "archived"] +Submission = Literal["UNSUBMITTED", "SUBMITTED"] SortField = Literal[ "id", "name", @@ -62,6 +63,7 @@ class PreferencesInput(Contract): class AlphaFilters(Contract): + submission: Submission | None = None q: str | None = Field(default=None, max_length=300) region: str | None = None universe: str | None = None @@ -160,16 +162,34 @@ class BulkUpdate(BulkInput): class JobInput(Contract): - kind: Literal["full_sync", "alpha_refresh", "pnl_refresh"] + kind: Literal["full_sync", "daily_sync", "alpha_refresh", "pnl_refresh", "self_correlation"] alpha_ids: list[str] = Field(default_factory=list) + submission: Submission | None = None + date_from: date | None = None + date_to: date | None = None @model_validator(mode="after") def validate_ids(self): - if self.kind == "full_sync": + if self.kind in ("full_sync", "daily_sync"): if self.alpha_ids: - raise ValueError("全量同步不接受 Alpha ID") + raise ValueError("列表同步不接受 Alpha ID") + if self.kind == "full_sync": + if self.submission == "UNSUBMITTED": + raise ValueError("待提交 Alpha 必须选择日期逐天同步") + self.submission = "SUBMITTED" + if self.date_from is not None or self.date_to is not None: + raise ValueError("全量同步不接受日期范围") + else: + if not self.submission or not self.date_from or not self.date_to: + raise ValueError("按天同步必须选择待提交/已提交和起止日期") + if self.date_from > self.date_to: + raise ValueError("开始日期不能晚于结束日期") + if self.date_to > datetime.now(timezone.utc).date(): + raise ValueError("同步日期不能晚于今天(UTC)") else: self.alpha_ids = valid_ids(self.alpha_ids) + if self.submission is not None or self.date_from is not None or self.date_to is not None: + raise ValueError("按 ID 操作不接受分组或日期范围") return self @@ -199,6 +219,7 @@ class AlphaSummary(BaseModel): date_submitted: datetime | None synced_at: datetime research: ResearchOutput + local_correlation: dict | None = None class AlphaDetail(AlphaSummary): @@ -230,6 +251,8 @@ class JobOutput(BaseModel): next_retry_at: datetime | None created_at: datetime updated_at: datetime + payload: dict = Field(default_factory=dict) + checkpoint: dict = Field(default_factory=dict) class PlatformSessionOutput(BaseModel): @@ -286,6 +309,11 @@ class PnlOutput(BaseModel): fetched_at: datetime | None +class SelfCorrelationOutput(BaseModel): + cached: bool + result: dict | None + + class FacetsOutput(BaseModel): region: list[str] universe: list[str] diff --git a/backend/app/worldquant.py b/backend/app/worldquant.py index e5aea35..e4d8eaf 100644 --- a/backend/app/worldquant.py +++ b/backend/app/worldquant.py @@ -257,20 +257,25 @@ class WqClient: errors[key] = str(exc) return usage_snapshot(data, errors) - async def alphas(self, submission, hidden, offset, before): + async def alphas(self, submission, hidden, offset, before, *, date_from=None, date_to=None): # Cover every platform stage; submitted records are not assumed to be OS only. status_key = "status" if submission == "UNSUBMITTED" else "status!" - return await self.get( - "/users/self/alphas", - { - status_key: "UNSUBMITTED", - "hidden": str(hidden).lower(), - "limit": 100, - "offset": offset, - "order": "dateCreated", - "dateCreated<": before, - }, - ) + params = { + status_key: "UNSUBMITTED", + "hidden": str(hidden).lower(), + "limit": 100, + "offset": offset, + "order": "dateCreated", + "dateCreated<": before, + } + if date_from is not None and date_to is not None: + # Daily intervals are [UTC midnight, next midnight). Submitted dates + # use the actual submission time, independently of creation or stage. + field = "dateCreated" if submission == "UNSUBMITTED" else "dateSubmitted" + params.update({f"{field}>=": date_from, f"{field}<": date_to, "order": field}) + if field == "dateCreated": + params["dateCreated<"] = min(before, date_to) + return await self.get("/users/self/alphas", params) async def alpha(self, alpha_id): return await self.get(f"/alphas/{alpha_id}") diff --git a/backend/migrations/versions/0003_local_self_correlation.py b/backend/migrations/versions/0003_local_self_correlation.py new file mode 100644 index 0000000..b9419fe --- /dev/null +++ b/backend/migrations/versions/0003_local_self_correlation.py @@ -0,0 +1,25 @@ +"""Persist local self-correlation independently of platform snapshots.""" + +from alembic import op +import sqlalchemy as sa + +revision = "0003" +down_revision = "0002" +branch_labels = None +depends_on = None + + +def upgrade(): + op.create_table( + "self_correlations", + sa.Column("alpha_id", sa.String(100), sa.ForeignKey("alphas.id"), primary_key=True), + sa.Column("region", sa.String(50), nullable=True), + sa.Column("result", sa.JSON(), nullable=False), + sa.Column("stale", sa.Boolean(), nullable=False), + sa.Column("calculated_at", sa.DateTime(timezone=True), nullable=False), + ) + op.create_index("ix_self_correlations_region", "self_correlations", ["region"]) + + +def downgrade(): + op.drop_table("self_correlations") diff --git a/backend/tests/browser_server.py b/backend/tests/browser_server.py index 15eb448..928bd27 100644 --- a/backend/tests/browser_server.py +++ b/backend/tests/browser_server.py @@ -45,7 +45,7 @@ def sample(index): "selection": {"code": "self_correlation < 0.5"} if super_alpha else None, "combo": {"code": "alpha"} if super_alpha else None, "settings": { - "region": ["USA", "CHN", "EUR"][index % 3], + "region": ["USA", "CHN", "EUR"][(index // 3) % 3], "universe": "TOP3000", "language": language, "delay": 1, @@ -64,7 +64,10 @@ def sample(index): "checks": [{"name": "LOW_SHARPE", "result": "PASS", "value": 2.1, "limit": 1.58}], }, "os": {"sharpe": 1.1} if index % 3 == 0 else None, - "dateCreated": (datetime(2025, 1, 1, tzinfo=timezone.utc) + timedelta(days=index)).isoformat(), + "dateCreated": (datetime(2025, 1, 1, tzinfo=timezone.utc) + timedelta(seconds=index)).isoformat(), + "dateSubmitted": (datetime(2025, 2, 1, tzinfo=timezone.utc) + timedelta(days=index % 2)).isoformat() + if index % 3 == 0 + else None, } @@ -141,6 +144,21 @@ def create_test_app(): matched = [ r for r in records if (r["status"] == "UNSUBMITTED") == unsubmitted and r["hidden"] == hidden ] + for field in ("dateCreated", "dateSubmitted"): + for suffix in (">=", "<"): + boundary = request.url.params.get(field + suffix) + if boundary: + bound = datetime.fromisoformat(boundary).replace(tzinfo=timezone.utc) + matched = [ + r + for r in matched + if r.get(field) + and ( + datetime.fromisoformat(r[field]) >= bound + if suffix == ">=" + else datetime.fromisoformat(r[field]) < bound + ) + ] offset, limit = ( int(request.url.params.get("offset", 0)), int(request.url.params.get("limit", 100)), diff --git a/backend/tests/test_alpha_management.py b/backend/tests/test_alpha_management.py new file mode 100644 index 0000000..e1705b5 --- /dev/null +++ b/backend/tests/test_alpha_management.py @@ -0,0 +1,341 @@ +"""Behavioral coverage for submission scopes, daily recovery and local correlation.""" + +import csv +import io +from datetime import date, datetime, timedelta + +import httpx +import pytest +from sqlalchemy import select + +from app.alphas import upsert_alpha +from app.correlation import calculate_correlation, daily_changes +from app.jobs import Runner +from app.models import Alpha, JobItem, Pnl, Research, SelfCorrelation +from app.worldquant import WqClient, WqError +from tests.conftest import alpha +from tests.test_jobs import FakePlatform, ready_runner, result + +PREFIX = "/api/v1" + + +def points(changes, start=date(2025, 1, 1), multiplier=1): + value, data = 100, [{"date": start.isoformat(), "value": 100}] + for index, change in enumerate(changes, 1): + value += change * multiplier + data.append({"date": (start + timedelta(days=index)).isoformat(), "value": value}) + return data + + +def reference(alpha_id, data): + return {"alpha_id": alpha_id, "points": data, "fetched_at": "2025-03-01T00:00:00Z"} + + +def test_signed_pearson_and_incomplete_coverage(): + changes = [((i * 7) % 19) - 8 for i in range(60)] + target = points(changes) + positive = reference("positive", points(changes, multiplier=2)) + negative = reference("negative", points(changes, multiplier=-1)) + computed = calculate_correlation(target, [positive, negative]) + assert computed["status"] == "high" and computed["compared_count"] == 2 + assert computed["most_correlated_alpha_id"] == "positive" + assert computed["matches"][0]["sample_count"] == 60 + assert computed["max_correlation"] == pytest.approx(1) + computed = calculate_correlation(target, [negative]) + assert computed["status"] == "low" and computed["max_correlation"] == pytest.approx(-1) + incomplete = calculate_correlation(target, [negative, {"alpha_id": "missing", "error": "无权访问"}]) + assert incomplete["status"] == "partial" and incomplete["skipped_count"] == 1 + + +@pytest.mark.parametrize("candidate", [[], points([1] * 60), points([1, 3, 2] * 9)]) +def test_insufficient_and_constant_series_are_never_passed(candidate): + report = calculate_correlation(points([1, 3, 2] * 20), [reference("invalid", candidate)]) + assert report["status"] == "insufficient_data" + assert report["max_correlation"] is None and report["compared_count"] == 0 + assert report["skipped_count"] == 1 + + +def test_missing_points_break_intervals_and_dates_align_before_comparison(): + original = points([1, 3, 2] * 20) + missing = [dict(p) for p in original] + missing[4]["value"] = None + changes, _ = daily_changes(missing) + assert date(2025, 1, 5) not in changes and date(2025, 1, 6) not in changes + report = calculate_correlation(original, [reference("gap", list(reversed(missing)))]) + assert report["matches"][0]["sample_count"] == 58 + # A missing row must not pair a two-day increment with a one-day increment. + removed = original[:4] + original[5:] + report = calculate_correlation(original, [reference("gap", removed)]) + assert report["matches"][0]["sample_count"] == 58 + with pytest.raises(ValueError, match="同一天"): + daily_changes(original + [original[0]]) + + +def test_common_four_year_window_is_anchored_to_target(): + changes = [1, 3, 2] * 20 + target = points(changes, date(2025, 1, 1)) + historical = points(changes, date(2019, 1, 1)) + report = calculate_correlation(target, [reference("old", historical)]) + assert report["status"] == "insufficient_data" and report["max_correlation"] is None + assert report["window_from"] == "2021-03-02" + + +async def test_submission_tabs_and_export_share_the_same_scope(app, logged_in): + async with app.state.sessions() as db: + for raw in ( + alpha("pending"), + alpha("active", status="ACTIVE", stage="IS"), + alpha("retired", status="DECOMMISSIONED"), + alpha("unknown", status=None), + ): + await upsert_alpha(db, raw) + await db.commit() + pending = (await logged_in.get(f"{PREFIX}/alphas?submission=UNSUBMITTED")).json() + submitted = (await logged_in.get(f"{PREFIX}/alphas?submission=SUBMITTED")).json() + assert [a["id"] for a in pending["items"]] == ["pending"] + assert {a["id"] for a in submitted["items"]} == {"active", "retired"} + exported = await logged_in.get(f"{PREFIX}/alphas/export?submission=SUBMITTED") + assert {r["id"] for r in csv.DictReader(io.StringIO(exported.text.lstrip("\ufeff")))} == { + "active", + "retired", + } + assert (await logged_in.get(f"{PREFIX}/alphas?submission=BAD")).status_code == 422 + + +async def test_daily_job_validation_scope_deduplication_and_history(app, logged_in): + await ready_runner(app) + for payload in ( + {"kind": "daily_sync"}, + {"kind": "full_sync", "submission": "UNSUBMITTED"}, + {"kind": "full_sync", "date_from": "2025-01-01"}, + {"kind": "daily_sync", "submission": "SUBMITTED", "date_from": "2025-01-02", "date_to": "2025-01-01"}, + {"kind": "daily_sync", "submission": "SUBMITTED", "date_from": "2025-01-01", "date_to": "9999-01-01"}, + {"kind": "alpha_refresh", "alpha_ids": ["a"], "submission": "SUBMITTED"}, + ): + assert (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).status_code == 422 + payload = { + "kind": "daily_sync", + "submission": "UNSUBMITTED", + "date_from": "2025-01-01", + "date_to": "2025-01-02", + } + first = (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).json() + duplicate = (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).json() + other = (await logged_in.post(f"{PREFIX}/sync-jobs", json={**payload, "submission": "SUBMITTED"})).json() + assert first["id"] == duplicate["id"] != other["id"] + assert first["payload"]["date_from"] == "2025-01-01" + full = (await logged_in.post(f"{PREFIX}/sync-jobs", json={"kind": "full_sync"})).json() + assert full["payload"]["submission"] == "SUBMITTED" + + +class DailyPlatform(FakePlatform): + def __init__(self): + super().__init__() + self.daily_calls = [] + self.fail_daily = True + self.records = [ + alpha("midnight", dateCreated="2025-01-01T00:00:00Z"), + alpha("end", dateCreated="2025-01-01T23:59:59.999999Z"), + alpha("next", dateCreated="2025-01-02T00:00:00Z"), + alpha("hidden1", hidden=True, dateCreated="2025-01-02T01:00:00Z"), + alpha("hidden2", hidden=True, dateCreated="2025-01-02T02:00:00Z"), + alpha( + "submitted", + status="ACTIVE", + dateCreated="2024-01-01T00:00:00Z", + dateSubmitted="2025-01-02T00:00:00Z", + ), + ] + + async def alphas(self, submission, hidden, offset, before, *, date_from=None, date_to=None): + self.daily_calls.append((submission, hidden, offset, date_from, date_to)) + if self.fail_daily and hidden and offset == 1: + raise WqError("模拟日内第二页失败", "network_error") + field = "dateCreated" if submission == "UNSUBMITTED" else "dateSubmitted" + matched = [ + r + for r in self.records + if (r["status"] == "UNSUBMITTED") == (submission == "UNSUBMITTED") and r["hidden"] == hidden + ] + if date_from: + matched = [ + r + for r in matched + if datetime.fromisoformat(date_from) + <= datetime.fromisoformat(r[field].replace("Z", "+00:00")) + < datetime.fromisoformat(date_to) + ] + return {"results": matched[offset : offset + 1], "count": len(matched)} + + +async def test_daily_sync_midnight_hidden_pages_restart_and_research_preservation(app, logged_in): + runner = await ready_runner(app) + runner.client = DailyPlatform() + payload = { + "kind": "daily_sync", + "submission": "UNSUBMITTED", + "date_from": "2025-01-01", + "date_to": "2025-01-02", + } + job_id = (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).json()["id"] + await runner.execute(job_id) + failed = await result(runner, job_id) + assert failed.status == "failed" and failed.processed == 4 + assert failed.checkpoint["date"] == "2025-01-02" and failed.checkpoint["offset"] == 1 + async with runner.sessions() as db: + research = await db.get(Research, "midnight") + research.note, research.state = "preserved", "candidate" + await db.commit() + resumed = Runner(runner.sessions, runner.settings, DailyPlatform()) + resumed.client.fail_daily = False + await resumed.execute(job_id) + complete = await result(resumed, job_id) + assert complete.status == "completed" and complete.processed == complete.total == 5 + assert complete.checkpoint["dates_completed"] == complete.checkpoint["dates_total"] == 2 + assert resumed.client.daily_calls[0][1:3] == (True, 1) + async with runner.sessions() as db: + assert len((await db.scalars(select(JobItem).where(JobItem.job_id == job_id))).all()) == 5 + assert (await db.get(Research, "midnight")).note == "preserved" + assert await db.get(Alpha, "submitted") is None + # Daily submitted sync uses submission date, although its creation was a year earlier. + submitted_id = ( + await logged_in.post(f"{PREFIX}/sync-jobs", json={**payload, "submission": "SUBMITTED"}) + ).json()["id"] + await resumed.execute(submitted_id) + assert (await result(resumed, submitted_id)).processed == 1 + full_id = (await logged_in.post(f"{PREFIX}/sync-jobs", json={"kind": "full_sync"})).json()["id"] + resumed.client.daily_calls.clear() + await resumed.execute(full_id) + assert (await result(resumed, full_id)).processed == 1 + assert {call[0] for call in resumed.client.daily_calls} == {"SUBMITTED"} + + +async def test_platform_ignoring_daily_filter_does_not_import_other_days(app, logged_in): + runner = await ready_runner(app) + + async def wrong_day(*args, **kwargs): + return {"results": [alpha("wrong", dateCreated="2024-01-01T00:00:00Z")], "next": None} + + runner.client.alphas = wrong_day + job = ( + await logged_in.post( + f"{PREFIX}/sync-jobs", + json={ + "kind": "daily_sync", + "submission": "UNSUBMITTED", + "date_from": "2025-01-01", + "date_to": "2025-01-01", + }, + ) + ).json() + await runner.execute(job["id"]) + assert (await result(runner, job["id"])).status == "failed" + async with runner.sessions() as db: + assert await db.get(Alpha, "wrong") is None + + +async def test_local_detection_uses_cache_excludes_self_and_keeps_research(app, logged_in): + runner = app.state.runner # No credentials: any upstream call fails this test. + data = points([1, 4, 2, -2] * 20) + async with runner.sessions() as db: + for raw in [ + alpha("target", status="ACTIVE"), + alpha("peer", status="ACTIVE"), + alpha("pending"), + alpha("other-region", status="ACTIVE", settings={"region": "CHN"}), + alpha("unknown", status=None), + ]: + await upsert_alpha(db, raw) + db.add(Pnl(alpha_id=raw["id"], raw={}, points=data)) + await db.flush() + research = await db.get(Research, "target") + research.note, research.state = "hypothesis", "candidate" + await db.commit() + job = ( + await logged_in.post( + f"{PREFIX}/sync-jobs", json={"kind": "self_correlation", "alpha_ids": ["target"]} + ) + ).json() + await runner.execute(job["id"]) + assert (await result(runner, job["id"])).status == "completed" + report = (await logged_in.get(f"{PREFIX}/alphas/target/self-correlation")).json()["result"] + assert report["candidate_count"] == report["compared_count"] == 1 + assert report["matches"][0]["alpha_id"] == "peer" and report["status"] == "high" + for timestamp in ( + report["calculated_at"], + report["target_pnl_fetched_at"], + report["matches"][0]["pnl_fetched_at"], + ): + assert datetime.fromisoformat(timestamp).utcoffset() == timedelta(0) + async with runner.sessions() as db: + assert (await db.get(Research, "target")).state == "candidate" + assert (await db.get(Research, "target")).note == "hypothesis" + assert (await db.get(Alpha, "target")).checks[0]["result"] == "FAIL" + # A new submitted reference invalidates the old result without deleting it. + await upsert_alpha(db, alpha("new-peer", status="ACTIVE")) + await db.commit() + assert (await db.get(SelfCorrelation, "target")).stale + summary = (await logged_in.get(f"{PREFIX}/alphas?q=target")).json()["items"][0]["local_correlation"] + assert summary["stale"] and summary["max_correlation"] == pytest.approx(1) + assert datetime.fromisoformat(summary["calculated_at"]).utcoffset() == timedelta(0) + + +async def test_missing_pnl_filled_once_and_partial_data_reported(app, logged_in): + runner = await ready_runner(app) + calls = [] + data = points([1, 3, -2] * 20) + + async def pnl(alpha_id): + calls.append(alpha_id) + if alpha_id == "unavailable": + raise WqError("无权访问", "access_denied") + multiplier = 1 if alpha_id == "target" else -1 + return {"records": [{"date": p["date"], "pnl": p["value"] * multiplier} for p in data]} + + runner.client.pnl = pnl + async with runner.sessions() as db: + for raw in [alpha("target"), alpha("peer", status="ACTIVE"), alpha("unavailable", status="ACTIVE")]: + await upsert_alpha(db, raw) + await db.commit() + + async def run(): + job = ( + await logged_in.post( + f"{PREFIX}/sync-jobs", json={"kind": "self_correlation", "alpha_ids": ["target"]} + ) + ).json() + await runner.execute(job["id"]) + return (await logged_in.get(f"{PREFIX}/alphas/target/self-correlation")).json()["result"] + + report = await run() + assert report["status"] == "partial" and report["skipped_count"] == 1 + await run() + assert calls.count("target") == calls.count("peer") == 1 + + +async def test_scoped_date_query_parameters_and_no_platform_check(settings): + requests = [] + + def handler(request): + requests.append(request) + return httpx.Response(200, json={"results": []}) + + client = WqClient(settings, transport=httpx.MockTransport(handler)) + client.credentials, client.authenticated = ("test@example.com", "test"), True + for submission in ("UNSUBMITTED", "SUBMITTED"): + await client.alphas( + submission, + True, + 100, + "2025-03-01T00:00:00+00:00", + date_from="2025-01-01T00:00:00+00:00", + date_to="2025-01-02T00:00:00+00:00", + ) + first, second = [dict(r.url.params) for r in requests] + assert first["dateCreated>="] == "2025-01-01T00:00:00+00:00" + assert first["dateCreated<"] == "2025-01-02T00:00:00+00:00" + assert second["dateSubmitted>="] == first["dateCreated>="] + assert second["dateSubmitted<"] == first["dateCreated<"] + assert "status!" in second and second["hidden"] == "true" + assert all(r.method == "GET" and r.url.path == "/users/self/alphas" for r in requests) + await client.close() diff --git a/docs/project-plan.md b/docs/project-plan.md index f839b9b..44c165b 100644 --- a/docs/project-plan.md +++ b/docs/project-plan.md @@ -30,10 +30,11 @@ - 主界面为侧栏、筛选区、可排序分页表格、列显隐、多选、详情抽屉;任务进度在独立面板查看。 - 列表以完整 flex 高度链填充可用区域,仅表体滚动,分页固定在底部;聊天展开后按剩余空间布局,保留编辑草稿。 -- 全量同步已提交/未提交、隐藏/可见 Alpha;指定 ID 导入及选中刷新;仅手动触发。 +- 列表分为待提交、已提交两个 Tab,分组作用于查询、选择及导出。待提交先选 UTC 创建日期范围后逐天同步;已提交按 UTC 提交日期逐天同步或全量同步,均覆盖隐藏/可见记录。支持指定 ID 导入及选中刷新;仅手动触发。 - 筛选 ID/名称/表达式、地区、Universe、类型、语言、平台状态、隐藏、日期、核心指标、本地标签/状态/收藏。 - 展示 REGULAR/SUPER、FASTEXPR/PYTHON、Selection/Combo、完整 settings、IS/OS、已有 checks。 - PnL 按需获取、缓存、曲线展示、刷新;可打开 BRAIN 原页面。 +- 本地自相关以已同步的同地区已提交 Alpha 为基准,排除自身;缓存缺失时只读补取 PnL。累计 PnL 作日变化,在目标最新日往前四年计算 Pearson,最少 30 个共同样本,带符号最大值达到 0.7 时告警。展示不足及跳过原因;结果单独保存,缓存或基准改变后标记待重算。支持单条及最多 100 条批量检测,不触发平台检查。 - 本地备注、标签、收藏与研究状态单独保存,同步不能覆盖。研究状态:inbox/candidate/optimizing/archived。 - 批量加减标签及改研究状态;CSV 按当前筛选与排序导出全部结果,不限制为 500 条。 - 指标缺失保留 null,不伪装成零;本地研究状态与平台状态、检查结果分开。 @@ -50,7 +51,7 @@ React + Semi Design + AI SDK UI 提供可调整宽度的聊天面板;FastAPI + 页面读取本地数据库。`/api/v1/auth` 管理登录,`/account` 管理配置与资料,`/alphas` 管理查询及研究记录,`/alphas/{id}/pnl` 读取缓存,`/sync-jobs` 创建、查询、取消和重试任务。 长任务返回 job ID;前端轮询。首期单后端进程运行异步任务,任务及分页检查点持久化。 每页原子落库、按 Alpha ID 更新、失败重试及重启恢复;429 遵守 Retry-After,其余暂时性错误有界退避。 -原始业务响应与结构化摘要分别保存,不存认证敏感字段。上游数据移动可能影响 offset 分页,通过 ID 去重及再次全量同步校正,不因一次未查到就删除本地记录。 +原始业务响应与结构化摘要分别保存,不存认证敏感字段。上游数据移动可能影响 offset 分页,通过 ID 去重及再次同步对应范围校正,不因一次未查到就删除本地记录。 ## 路线图 @@ -59,7 +60,7 @@ React + Semi Design + AI SDK UI 提供可调整宽度的聊天面板;FastAPI + | 一 | 账户、列表、研究记录、同步、部署 | 本文件首期范围 | | 一扩展 | AI 聊天、模型配置、只读业务工具、修改确认闭环 | [AI Chatbot 开发计划](ai-chatbot-plan.md) | | 二 | 数据集/字段/算子、模板、批次队列、AST 校验、实验去重、暂停恢复 | 旧系统采样→密度→深度回测,以及 [回测台账](https://mail.google.com/mail/#all/19ea68a7dde5ceaa) | -| 三 | PnL 稳定性、比较、相关性、稳健性、跨区变体、Super Alpha 组合 | 旧系统有效分析能力 | +| 三 | PnL 稳定性、比较、进一步的相关性与稳健性分析、跨区变体、Super Alpha 组合 | 基础本地自相关已纳入当前 Alpha 管理;其余参考旧系统有效分析能力 | | 四 | 假设与实验记录、CLI/MCP、论坛检索与进一步的研究编排 | [决策摘要](https://mail.google.com/mail/#all/19fc16ce3f17311c)、[可复盘流程](https://mail.google.com/mail/#all/1a00ee69df671d4c) | | 后续 | 平台回写、检查、提交、顾问表现 | 另行确认业务范围 | diff --git a/frontend/src/ai/ChatPanel.tsx b/frontend/src/ai/ChatPanel.tsx index 83de847..5046a7e 100644 --- a/frontend/src/ai/ChatPanel.tsx +++ b/frontend/src/ai/ChatPanel.tsx @@ -579,7 +579,21 @@ function BusinessCard({ {Array.isArray(call.preview.operation.alpha_ids) && call.preview.operation.alpha_ids.length ? call.preview.operation.alpha_ids.join("、") - : "全部 Alpha"} + : call.preview.operation.submission === "UNSUBMITTED" + ? "待提交 Alpha" + : call.preview.operation.submission === "SUBMITTED" + ? "已提交 Alpha" + : "全部 Alpha"} + {call.preview.operation.kind === "daily_sync" && ( + <> + {" · "} + {call.preview.operation.submission === "UNSUBMITTED" + ? "创建日期" + : "提交日期"}{" "} + {String(call.preview.operation.date_from)} 至{" "} + {String(call.preview.operation.date_to)}(UTC) + > + )}
)} {call.preview.job && ( diff --git a/frontend/src/api.ts b/frontend/src/api.ts index 6bad5b6..0febade 100644 --- a/frontend/src/api.ts +++ b/frontend/src/api.ts @@ -83,12 +83,20 @@ export const stateOptions = Object.entries(stateLabels).map( ); export const jobLabels: Record
diff --git a/frontend/src/components/AlphaSyncDialog.tsx b/frontend/src/components/AlphaSyncDialog.tsx
new file mode 100644
index 0000000..d5f5875
--- /dev/null
+++ b/frontend/src/components/AlphaSyncDialog.tsx
@@ -0,0 +1,98 @@
+import { useState } from "react";
+import { Input, Modal, Toast } from "@douyinfe/semi-ui-19";
+import { post } from "../api";
+import type { Submission } from "../types";
+
+export function AlphaSyncDialog({
+ submission,
+ suspended,
+ onClose,
+ onTask,
+}: {
+ submission: Submission;
+ suspended: boolean;
+ onClose: () => void;
+ onTask: () => void;
+}) {
+ const [from, setFrom] = useState("");
+ const [to, setTo] = useState("");
+ const [busy, setBusy] = useState(false);
+ const [error, setError] = useState("");
+ const label = submission === "UNSUBMITTED" ? "待提交" : "已提交";
+ const dateLabel = submission === "UNSUBMITTED" ? "创建日期" : "提交日期";
+ const today = new Date().toISOString().slice(0, 10);
+ async function start() {
+ if (!from || !to || from > to || to > today) {
+ setError("请选择有效的起止日期,结束日期不能晚于今天(UTC)。");
+ return;
+ }
+ setBusy(true);
+ try {
+ await post("/sync-jobs", {
+ kind: "daily_sync",
+ submission,
+ date_from: from,
+ date_to: to,
+ });
+ onClose();
+ onTask();
+ } catch (e) {
+ Toast.error((e as Error).message);
+ } finally {
+ setBusy(false);
+ }
+ }
+ return (
+
+ 按{dateLabel}(UTC)逐天获取,包含隐藏记录。相同起止日期表示只同步当天。
+
+ {error}
+
+ 任务面板显示当前日期和进度,中断后可继续未完成部分。
+
{formatTime(job.created_at, timezone)}
+ {job.payload?.submission && ( ++ {job.payload.submission === "SUBMITTED" ? "已提交" : "待提交"} + {job.payload.date_from + ? ` · ${job.payload.date_from} 至 ${job.payload.date_to}(UTC)` + : " · 全量"} +
+ )} + {job.checkpoint?.date && ( ++ 同步日期:{job.checkpoint.date} · 已完成{" "} + {job.checkpoint.dates_completed} / {job.checkpoint.dates_total} 天 +
+ )} + {job.kind === "self_correlation" && job.checkpoint?.alpha_id && ( ++ {job.checkpoint.alpha_id} ·{" "} + {job.checkpoint.phase === "calculating" + ? "计算相关性" + : `准备 PnL:${job.checkpoint.references_loaded ?? 0} / ${job.checkpoint.references_total ?? 0} 个基准`} +
+ )}+ 与本地已同步的同地区已提交 Alpha 比较,排除自身。使用 PnL + 缓存,缺失时自动获取;建议先全量同步已提交 Alpha。 +
++ 近四年 PnL 日变化的 Pearson 相关系数,至少 30 个共同样本;本地告警线为 + 0.7,检测结果与平台检查分别记录。 +
+ {error &&{result.reason}
} + {result.status === "partial" && ( +已有样本低于阈值,但部分基准无法比较,不能据此判断全部样本。
+ )} + {result.matches.length > 0 && ( + <> +