diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..dcd12d8 --- /dev/null +++ b/Makefile @@ -0,0 +1,89 @@ +# tradingSystem (PMS) 常用操作一句话入口 +# ============================================================================ +# 为什么要这个文件: 日常命令都是「docker compose run --rm --no-deps pms-web python +# scripts/xxx.py」这种一行八十字符的东西, 手打容易漏 --no-deps (于是顺带把 beat/worker +# 也拉起来)、漏 --rm (于是攒一堆退出的容器)。把口径固定在这里, 少一类手误。 +# +# make help 看有哪些目标 +# make deploy 拉代码 → 重建镜像 → 重建容器 → 建表 (= scripts/deploy.sh) +# make test 跑全部单测 (不连库, 秒级) +# +# **改了 Python 代码必须 build + force-recreate**, 因为源码是打进镜像的 —— +# `docker compose restart` 跑的还是旧镜像里的旧代码, 而且一点报错都没有 (deploy.sh +# 头部有详细说明)。`make up` 用的就是这条路径, 别用 `docker compose restart` 代替。 +# ============================================================================ +.DEFAULT_GOAL := help +SHELL := /bin/bash + +# 各 profile 的组合。只想跑页面: `make up PROFILES=` +PROFILES ?= --profile sched --profile ws +DC := docker compose $(PROFILES) +# 一次性命令一律 --no-deps: 不然跑个单测都会把 beat/worker 拉起来 +RUN := docker compose run --rm --no-deps pms-web + +.PHONY: help deploy deploy-local build up down ps logs test initdb check health \ + probe changes industry ws-status reset-ledger shell + +help: ## 列出所有目标 + @grep -hE '^[a-zA-Z_-]+:.*?## .*$$' $(MAKEFILE_LIST) \ + | awk 'BEGIN {FS = ":.*?## "}; {printf " \033[36m%-14s\033[0m %s\n", $$1, $$2}' + +# ---------------------------------------------------------------- 部署 +deploy: ## 一句话部署: git pull → build → up --force-recreate → 建表 + @./scripts/deploy.sh + +deploy-local: ## 同 deploy 但跳过 git pull (本地已改好) + @./scripts/deploy.sh --no-pull + +build: ## 只重建镜像 + $(DC) build + +up: ## 起服务 (force-recreate —— 让新镜像真正生效) + $(DC) up -d --force-recreate + @$(DC) ps + +down: ## 停掉全部服务 (不删卷) + $(DC) down + +ps: ## 看服务状态 + @$(DC) ps + +logs: ## 跟日志 (make logs S=pms-ws 只看一个服务) + $(DC) logs -f $(S) + +# ---------------------------------------------------------------- 自检 +test: ## 全部单测 (零外部依赖, 不连库; 应输出 ALL SUITES PASS) + $(RUN) python scripts/run_tests.py + +initdb: ## 建表/补表 (幂等; 加 DRY=1 只演练打印) + $(RUN) python scripts/init_db.py $(if $(DRY),,--yes) + +check: ## 实机连通性与表结构自检 (需真实 .env) + $(RUN) python scripts/check_db.py + +health: ## 页面健康检查 (配置装载 + 库连通自证) + @curl -s http://127.0.0.1:38100/health | python3 -m json.tool + +# ---------------------------------------------------------------- 上游 / 行业 / 通道 +probe: ## 上游计划实机探活 (只读; 加 SNAP=1 顺带落一份名册快照) + $(RUN) python scripts/probe_plan_api.py $(if $(SNAP),--snapshot,) + +changes: ## 榜单变化: 谁新进谁掉榜、持仓票有没有被上游摘了 + @curl -s 'http://127.0.0.1:38100/api/upstream/plan-changes?limit=30' \ + | python3 -m json.tool + +industry: ## 行业源实探 (ready=false 就是硬拦截失效, 别当它绿着) + @curl -s http://127.0.0.1:38100/api/industry | python3 -m json.tool + +ws-status: ## ws 通道状态 (连接态 / seq 水位 / 出口队列) + $(RUN) python scripts/ws_smoke.py status + +# ---------------------------------------------------------------- 危险操作 (要确认) +reset-ledger: ## 清空账本重来 (影子期专用; 必须 CONFIRM=1) + @if [ "$(CONFIRM)" != "1" ]; then \ + echo "这会清空 PMS 账本。确认请加 CONFIRM=1, 例如:"; \ + echo " make reset-ledger CONFIRM=1 ARGS='--purge-channel --reset-ws'"; exit 1; fi + $(RUN) python scripts/reset_ledger.py --yes $(ARGS) + +shell: ## 进容器 (排查用) + docker compose run --rm --no-deps pms-web bash diff --git a/README.md b/README.md index ac4eede..50750a0 100644 --- a/README.md +++ b/README.md @@ -13,8 +13,9 @@ | `POSITION_MGMT_DESIGN.md` | 总体设计 **V0.4(定稿,开发启动)**:命令系统与管理页面/账本/仓位框架/动作引擎/两道关口/择时执行/下游通道。功能一次性开发,上线按依赖分三步切换 | | `QMT_WS_PROTOCOL.md` | **PMS ↔ QMT WebSocket 指令与回报协议 V1.0(定稿)**:传输与重连、Ed25519 签名与幂等、消息集、状态机、断线补发与对账兜底、部署前检查清单。**这是下发通道的唯一实现依据** | | `QMT_INTERFACE_REQUIREMENTS.md` | 与 QMT 侧的数据与接口需求清单 **V2.0**:A 部分(只读数据)与 C 部分(切换约定)有效;**B 部分的表通道已废止**,改由上面的 ws 协议承担 | -| `UPSTREAM_PLAN_API.md` | **上游选股计划接口 (`/plan`) 的接入记录**:应答结构、PMS 侧五条口径、参数表,以及**待上游确认的 10 个口径问题**(`upside` 单位、`date` 语义、分页、`changes` 结构、`theme` 稳定性…)| -| `ddl_pms_v1.sql` | PMS 全部自有表建表语句(153 代理侧,**14 张**:设计 §11 的 10 张 + ws 通道 3 张 + 现金流水 1 张) | +| `UPSTREAM_PLAN_API.md` | **上游选股计划接口 (`/plan`) 的接入记录**:应答结构、PMS 侧五条口径、参数表,**行业源换 gp_hybk 的定案(§8)**,以及**榜单变化改由 PMS 自算的口径与取舍(§9)**| +| `ddl_pms_v1.sql` | PMS 全部自有表建表语句(153 代理侧,**15 张**:设计 §11 的 10 张 + ws 通道 3 张 + 现金流水 1 张 + 计划名册快照 1 张) | +| `Makefile` | 常用操作一句话入口(`make help` 看全部)。改了代码用 `make deploy` / `make up`,**别用 `docker compose restart`**——源码打进镜像,restart 跑的还是旧代码且一声不吭 | | `config/settings.py` | 配置(基础设施键名对齐 bionic;业务参数为初值,页面调参持久化到 `pms_runtime_param` 后优先) | ## 模块地图 @@ -31,10 +32,11 @@ app/ exec_timing.py 择时实现B: 分日配额 / 分笔 / 买卖出手判定 / 14:45 兜底 / 窗口收口 action_engine.py 动作引擎: FILL 回踩补足 / ADD 盈利加仓 / DCA 补仓 / TRIM 保垫减仓 signal_rules.py 决策系统两条信号流的解析与消化口径 (含置信度尺度归一) + plan_diff.py 上游榜单的版本比对: 名册指纹 / 新进掉榜 / 档位升降 / 榜尾噪音闸 tradedays.py 交易日历: 调度守卫与执行窗口计算 ws_codec.py QMT 协议编解码: 规范化串 / Ed25519 签名验签 / 信封 / seq 水位推进 db/session.py 三库连接 + **严格单表访问守卫** (JOIN/逗号连表/跨表子查询一律拒绝) - repo/ 单表数据访问: pms_repo (自有 10 表) / qmt_repo (ws 通道 3 表) + repo/ 单表数据访问: pms_repo (自有 12 表) / qmt_repo (ws 通道 3 表) / downstream_repo (下游只读三表) services/ 编排层 param_store.py 运行参数中心 (表值优先于 settings 初值, 页面调参即时生效) @@ -47,8 +49,9 @@ app/ signal_service.py 盘中信号订阅 (db2 广播 + db3 风控卖出) → 卖出指令或提议 ledger_service.py 成交回放 / 对账 / 除权 / 盘前 / 日终结算 / 运营日报 market.py 行情 (Redis db13) 与参考位 (决策系统主口径 + 兜底自算) - plan_feed.py **上游选股计划 /plan 接入**: 解析 / 交易日龄校验 / 候选筛选 / theme 灌行业映射 - industry.py 行业划分可插拔适配器 (custom_table / gp_stock_category / 停用) + plan_feed.py **上游选股计划 /plan 接入**: 解析 / 交易日龄校验 / 候选筛选 + / 名册快照与榜单变化 (上游 changes 恒为 null, PMS 自己算) + industry.py 行业划分可插拔适配器 (gp_hybk / custom_table / 停用; 带实探不靠参数) ws/runner.py **常驻连接进程 (pms-ws)**: 握手/心跳/重连/补发 + 出口出栈 + 上行落库确认 web/ FastAPI + 单页 (Vue3 + ElementPlus),页面四块 + 运维/日报抽屉 scheduler.py Celery beat 调度总表 (九个调度位 + 三条守卫) @@ -61,12 +64,14 @@ scripts/ test_batch5_units.py 决策系统信号流解析与消化口径 8 例 test_batch6_units.py ws 通道: 测试向量/签名/公钥/水位/DDL/逐笔入账 65 例 test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 32 例 + test_batch8_units.py 榜单变化: 名册指纹/三种语义/尾部闸/落库往返 60 例 test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 58 例 init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练; 含 DDL 体检) check_db.py 实机连通性与表结构自检 (需真实 .env) gen_keys.py ws 通道密钥: 生成 / 只取公钥(--pubkey) / PEM(--pem) / 自检(--check) reset_ledger.py 清空账本并把回放游标对齐到当前 (影子运行期重来一次; 不碰下游表) probe_plan_api.py 上游计划实机探活: 通不通 / 字段口径 / 候选筛选结果 / 有没有价 + / 榜单变化 (只读; --snapshot 才落库, 那是它唯一的写操作) ws_smoke.py ws 联调工具: status/watch/place/cancel/inbox (绕开 dispatch_mode) ``` @@ -80,6 +85,8 @@ scripts/ 服务共用一个镜像:`pms-web`(管理页面,端口 38100)+ `pms-beat` / `pms-worker`(Celery 调度与执行,挂在 `sched` profile 下)。 +日常操作有 `Makefile` 包了一层(`make help` 看全部):`make deploy` 一句话完成拉代码→重建镜像→重建容器→建表;`make test` / `make check` / `make probe` / `make changes` / `make industry` / `make ws-status` 是自检与实探。下面是这些命令展开后的样子——不用 make 也照样能跑,两者等价。 + ```bash # 服务器首次部署 git clone <仓库地址> && cd tradingSystem @@ -90,7 +97,7 @@ docker compose build # 默认走清华 PyPI 镜像; 可 -- docker compose run --rm pms-web python scripts/run_tests.py # 建表 (幂等; 不加 --yes 只演练打印) docker compose run --rm pms-web python scripts/init_db.py --yes -# 实机自检 (连库, 需 .env): 库连通 + pms_* 十表 + 下游表完整列定义 + 行情 Redis +# 实机自检 (连库, 需 .env): 库连通 + pms_* 十五表 + 下游表完整列定义 + 行情 Redis docker compose run --rm pms-web python scripts/check_db.py docker compose up -d # 管理页面 @@ -103,10 +110,12 @@ docker compose logs -f pms-beat pms-worker docker compose --profile ws up -d pms-ws # 启用 QMT 直连 (先配好 .env 里的两把密钥) docker compose logs -f pms-ws -# 日常更新 -git pull && docker compose build && docker compose up -d +# 日常更新 (= make deploy) +./scripts/deploy.sh ``` +**改了 Python 代码必须 build + force-recreate**:源码是打进镜像的,`docker compose restart` 只是把老容器停了再起,跑的还是旧镜像里的旧代码——**而且一点报错都没有**。`make deploy` / `make up` 走的都是正确路径。 + 基础镜像 `python:3.11-slim` 拉取慢时,先给服务器 Docker 配置 registry 镜像加速。日志落 `./logs`(已挂载卷);容器时区 Asia/Shanghai。 **管理页面的两个坑**(都踩过,写在这里): @@ -197,9 +206,17 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ 两条无条件覆盖档位的规矩:**减持方向不设确认门槛**(TRIM 保垫减仓任何档位都直接落指令);**−15% 及更深的补仓永远需用户确认**(即便档位是 full 也强制入队)。同一只票的同一动作若已有在途提议或在途指令,不重复提。 +## 榜单变化提示(只提示,不产生指令) + +上游 `/plan` 应答里的 `changes` 字段至今恒为 `null`,所以这件事改成 PMS 自己算:每次拉计划落一份名册快照(`pms_plan_snapshot`),跟上一份比出新进榜 / 掉榜 / 档位升降 / 券商覆盖翻转 / 名次跳变。页面在「上游计划」抽屉里,**持仓票单独一段**——上游把一只持仓票摘出榜或降了档,是「还该不该继续拿着」的直接信号,混在几十条新进榜里会被淹掉。 + +**它和上面的自主提议不是一回事,别混。** 自主提议会落指令或入确认队列;榜单变化**只是提示**,从头到尾不碰账本、不进规则闸、不产生任何指令。要减要清仍走命令台或提议确认。 + +两个容易踩的点,都已经处理:**掉榜不能直接算「上一版有、这一版没有」**——榜是按 `top=300` 截断过的,从 rank 250 掉到 310 只是被截掉,不是上游摘了它;新进榜是对称的同一个问题。**同一个 `date` 的计划会变**(`/plan` 是实时算的,应答里没有版本戳),所以「同一计划日的新版本」单独算一类并挂黄色横幅——它意味着此前拿到的候选池与现在不是同一份。口径与取舍详见 `UPSTREAM_PLAN_API.md` §9。 + ## 已实现 / 待开发 -**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 244 例**。 +**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**榜单变化提示**(名册快照 + 新进/掉榜/档位升降/覆盖翻转/名次跳变,持仓票单列,榜尾截断噪音闸);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 304 例**。 ### 下一步(按可动工顺序) @@ -211,12 +228,12 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ | 4 | **账本重建**:清账后账本为空,需对端持仓就绪后 `POST /api/ops/reconcile?apply_fix=true` 以下游为准补回 | 阻塞:等 QMT 侧按真实成本价重建模拟持仓(`QMT_SIDE_S3_CLOSEOUT.md` §5) | | 5 | **ws 通道联调收尾**(协议 §9 S3) | 阻塞:`trade_no` 格式不合 §5.5(`QMT_SIDE_S3_CLOSEOUT.md` §1,这条挡住切 ws) | | 6 | 上游计划的剩余待确认口径 | 等上游:`UPSTREAM_PLAN_API.md` §4 剩 5 条 + §7.3 新增两条(同一 `date` 计划不幂等、档位规模与文档不符) | -| 7 | `changes` 段做成页面提示(新进传导链 / 掉榜) | 可做,上游已确认该段是升降档,每日几十只 | -| 8 | T0 做T(二期) | 可做,设计已有,无外部依赖 | +| 7 | ~~榜单变化做成页面提示(新进传导链 / 掉榜)~~ | ✅ 2026-07-31:上游 `changes` 恒为 `null`,所以**改成 PMS 自己算**——每次拉计划落一份名册快照(第 15 张表),差异由纯逻辑比。顺带补上「同一 `date` 多版本、下游无从分辨手上是哪一版」那个洞。详见 `UPSTREAM_PLAN_API.md` §9 | +| 8 | T0 做T(二期) | 可做,设计已有,无外部依赖。**建议排在账本重建之后**:它要拿真实持仓验触发与 14:50 平回,账本空着只能跑单测 | | 9 | 择时实现 A(委托决策系统盘中择时) | 阻塞:等 bionic 侧接口 | | 10 | 研判闸接通 | 阻塞:等 bionic 侧 `process_intraday_audit` 新增 PMS direction。客户端已就位,接口好了在页面填 `PMS_JUDGE_API_BASE` 即通 | | 11 | `trading_buy_plan` 退场:PMS 已不读它 | 待上游确认无其他消费方(`UPSTREAM_PLAN_API.md` Q10) | -| 12 | 部署便利性:`make deploy` 一句话完成 build + up + 各 profile | 待办 | +| 12 | ~~部署便利性:`make deploy` 一句话完成 build + up + 各 profile~~ | ✅ 2026-07-31:`Makefile`(`make help` 看全部)。`deploy` 包 `scripts/deploy.sh`;另有 `test` / `initdb` / `check` / `probe` / `changes` / `industry` / `ws-status`。一次性命令统一带 `--no-deps`,免得跑个单测把 beat/worker 也拉起来 | ### 运行态注意(2026-07-31 收尾时的状态) @@ -224,6 +241,7 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ - 账本已清空(`reset_ledger --purge-channel --reset-ws --drop-reports`),ws seq 水位归零 → pms-ws 一重连对端会从 seq 1 全量补发约 1.2 万条。账本空的时候**不要**点「成交回放」。 - 宿主机放行 docker 网段到 8300 的规则若用 `iptables -I` 加的,**重启就没了**,需 `netfilter-persistent save` 或改 firewalld permanent。不持久化的话服务器重启后候选池会静默变空。 - 自主档位仍是 `propose_only`,下发通道仍是 `shadow`。 +- 新增第 15 张表 `pms_plan_snapshot`(名册快照),**部署后要跑一次 `make initdb`**,否则榜单变化那段会一直报「读名册快照失败:表不存在」——候选池不受影响。第一份快照落下之前不报任何变化(`kind=first`),这是设计如此不是坏了。 ### ws 通道实现清单 diff --git a/UPSTREAM_PLAN_API.md b/UPSTREAM_PLAN_API.md index 06cd5e1..9f447e0 100644 --- a/UPSTREAM_PLAN_API.md +++ b/UPSTREAM_PLAN_API.md @@ -528,3 +528,137 @@ curl -s http://127.0.0.1:38100/api/industry | python3 -m json.tool > 提醒: 行业源通了以后 `PMS_SECTOR_MAX_NAMES=4` / `PMS_SECTOR_MAX_RATIO=40%` 才真正开始 > 拦人。参考项目在**三级**上用的阈值是 20%。40% 是当初按"二级或更粗"的粒度定的, 三级粒度下 > 偏松 —— 等账本重建、有真实持仓分布之后再回头看这个数, 现在不动。 + + +--- + +## 9. 榜单变化改成 PMS 自算 (2026-07-31) + +上游应答里的 `changes` 字段, **至今回的都是 `null`** (§4 Q4 已问, 只拿到语义答复 +「升降档: 覆盖新增、估值闸翻转、传导链进出, 每日几十只」, 没给过非空示例)。而 +「谁新进榜、谁掉榜、谁降了档」是上游观点变化最直接的信号 —— 尤其**持仓票被摘出榜或降了 +档**, 那是"还该不该继续拿着"的直接依据。 + +等对方补, 就是把一件自己能算的事挂在外部依赖上。所以改成: **PMS 每次拉计划落一份名册 +快照, 差异自己算。** + +### 9.1 一石二鸟: 顺带补上 §7.3 那个洞 + +§7.3 记的是「同一 `date` 的计划会变」: `/plan` 不是读静态文件, 是 `plan.collect()` 实时 +从库里算的, 07-31 实测相隔一小时的两次请求给出了不同结果 (打分池 972/81 → 423/502、 +榜首换人)。当时提的改进是「请上游加 `generated_at` / 版本戳」, 因为**下游没有任何办法 +判断"手上这份是哪一版"**。 + +落了快照, 这个问题就地解决, 一个字都不用上游改: 同一 `plan_date` 底下有几行快照, 就是 +上游那天重算过几版; 每一行都带指纹与拉到的时刻。复盘时能确切说出"我当时用的是这一版", +而不是靠猜。 + +### 9.2 落点 + +| 环节 | 落点 | +|---|---| +| 表 | `pms_plan_snapshot` (**第 15 张**, DDL 尾部)。名册存 JSON, 一份 300 只约 20KB | +| 纯逻辑 | `app/core/plan_diff.py` —— 名册归一 / 指纹 / 比对。零外部依赖 | +| 服务层 | `plan_feed.snapshot()` 落库 · `changes()` 页面用 · `preview_changes()` 只读 | +| 接口 | `GET /api/upstream/plan-changes` | +| 页面 | 「上游计划」抽屉里的「榜单变化」段: 持仓票 / 掉榜与降档 / 新进榜 / 升档 / 名次跳变 / 快照台账 | +| 探活 | `probe_plan_api.py` 的 `[7]` 段 (**默认只读**, `--snapshot` 才写) | +| 单测 | `scripts/test_batch8_units.py` 60 例 | + +参数 (页面可改): + +``` +PMS_PLAN_SNAPSHOT true 每次拉计划落一份名册快照 +PMS_PLAN_SNAPSHOT_KEEP 200 保留份数 +PMS_PLAN_DIFF_RANK_JUMP 50 名次跳变多少名才报 +PMS_PLAN_DIFF_TAIL_GUARD 0.5 榜尾进出的噪音闸, 见 9.4 +``` + +### 9.3 三种比对语义要分开 + +混在一起看会得出错误结论: + +| kind | 含义 | 页面怎么显示 | +|---|---|---| +| `first` | 没有上一份快照 | **一条变化都不报** —— 否则首次是整整一榜 300 条"新进榜" | +| `cross_day` | 计划日不同 | Q4 说的「每日几十只」的正常升降档 | +| `same_date_revision` | **同一计划日的不同版本** | 黄色横幅。这不是市场变化, 是上游重算 —— 它意味着此前拿到的候选池与现在不是一回事 | + +### 9.4 掉榜不能直接算「上一版有、这一版没有」 + +我们向上游要 `top=300`。一只票从 rank 250 掉到 rank 310, 会从名册里"消失"—— +但那不是上游摘了它, 只是排到了我们要的条数之外。**吃满 top 的那一版, 尾部的进出全是 +截断噪音。** 不作区分的话每天多出几十上百条假掉榜, 提示就没人看了 —— 跟 §5 里 +`_capped` 治的那个「拿 counts.main=961 去比返回条数」是同一类错误的两个面。 + +新进榜是**对称**的同一个问题。所以两道闸各看各的那一版: + +``` +掉榜 —— 看 这一版 吃没吃满 (它截断了, 掉出去的可能只是被截掉) +新进 —— 看 上一版 吃没吃满 (它截断了, 新面孔可能上一版就在, 只是没露面) +``` + +闸位是 `PMS_PLAN_DIFF_TAIL_GUARD` (默认 0.5 = 排名进入榜单前一半才算数), 被挡下的计进 +`tail_churn` 只报个数、不列名。**没吃满 top 的那一版不设闸** —— 那种情况下的进出是真的。 +另有一条: 某一档在这一版一条都没有时也不设闸, 否则"整档消失"这种大事会被闸吞掉。 + +这道闸有三条边界, 每一条都是"别把真信号当噪音吞了": + +1. **持仓票完全豁免。** 闸是拿"少报几条噪音"换"提示还有人看", 这笔账在持仓票上不成立: + 非持仓票误报一条只是白看一眼, **持仓票漏报一条就是一个该减没减的仓位**。按默认 + `top=300` + `tail_guard=0.5`, 不豁免的话主榜后一半是个永久盲区。豁免出来的条数记在 + `tail_churn.held_exempt`, 页面明写。 +2. **两档各判各的。** `capped` 按档给 —— `top` 与 `obs_top` 是两个独立的请求参数, 主榜吃满 + 完全不代表观察档也吃满。混用一个标志的后果是: 观察档一次真实的摘牌被记成"截断噪音"。 +3. **`rank` 判不出来时照报。** 本模块别处的口径是"判不了就不判"= 不报变化; 但这个判据的 + "是"代表**吞掉一条变化**, 所以方向要反过来 —— 拿不准就放行。 + +### 9.5 几个刻意的取舍 + +- **指纹不含 `score`。** `/plan` 实时算, score 末位天天抖; 算进指纹, 每次缓存过期重拉都会 + 多落一行快照, 表白涨而信息量为零。含 `rank` 就够 —— 分数抖到改变了次序才算真的变了, + 没改变次序的抖动对候选池毫无影响 (候选就是按 score 降序切前 N)。 +- **未知档位不参与升降判定。** 档位强弱按 Q7 的「强传导 > 弱传导 > 无传导」; 上游哪天加了 + 新档位, 宁可不报, 不能报反。 +- **主榜 ↔ 观察档 互换单列一段**, 不算掉榜。观察档的定义是"无券商覆盖", 所以这个变化的 + 含义是**券商覆盖翻转**(Q4 说的「估值闸翻转」), 票还在, 只是没有估值锚了。§7.3 那次 + 423/502 的大对调就是这一类, 混进"掉榜"里会完全看不出发生了什么。 + 但**持仓票掉进观察档要算进 `n_watch`**(页面红条的判据): 观察档的行没有 `tier`, 所以它 + 既不算掉榜也不算降档, 不显式算进来的话, 一只持仓票丢了估值锚而页面一片安静。 +- **`digest` 上不设唯一键, 去重只跟"最新那行"比。** 当初想拿 `(plan_date, digest)` 唯一键 + 当去重手段, 有两个坑: ① SQLAlchemy 的 MySQL 方言默认开 `CLIENT_FOUND_ROWS`, + `ON DUPLICATE KEY UPDATE` 命中重复照样回 1 影响行, 拿 rowcount 判"是不是新插的"必错 + (`inbox_put` 那里已经栽过一次); ② 更要命的是同一天上游 **A→B→A 改回去**时, 第三次的 A + 因为"曾经存过"而落不下去, 库里最新仍是 B —— 页面会把 B→A 这次**真实回退**显示成 A→B, + 方向整个反了, 还顺带诬告一句"手上这份没落进快照 (落库失败?)"。现在改成先查最新一行、 + 变了才插。代价是并发下可能多落一行相同的快照 (比对结果为"无变化", 不产生假信号)。 +- **上一版"存在但读不出来"要说成故障。** `_loads` 解析失败返回 `[]`, 长得跟"没有上一版" + 一模一样 —— 于是一条真实变化都不报, 页面还显示得很正常。现在 `prev_broken` 单独标出来, + 页面挂红条明说"这次比不了, 空白不代表没有变化"。 +- **换了档就不算名次跳变** —— 两档的 rank 各排各的, 不可比。 +- **只提示, 不产生任何指令。** 持仓票掉榜是"要不要继续拿着"的信号, 但减仓/清仓仍走命令台 + 或提议确认。这条通道从头到尾不碰账本、不进规则闸。 +- **落库失败不阻断候选池**, 但要把错带出来 (`_snapshot_quiet` 的口径, 同 theme 灌库)。 + 另有 `in_sync`: 手上这份计划的指纹与库里最新快照对不上时, 页面直接说「下面比的是两个旧 + 版本」—— 宁可摆出来, 也不显示一个看起来正常、实际过时的差异。 +- **`reset_ledger` 不清这张表。** 它记的是**上游给过我们什么**, 与本方账本无关; 清账是 + "我方重来一次", 不该把上游的历史一起抹掉。 + +### 9.6 实机验收 + +```bash +# 【服务器 factorevaluation · ~/project/tradingSystem】 +make test # 应输出 ALL SUITES PASS (304 例) +make initdb # 幂等, 建第 15 张表 +make check # [2] 段应能看到 pms_plan_snapshot +make probe SNAP=1 # 落第一份快照; [7] 段此时应报"库里还没有快照" +make probe SNAP=1 # 再跑一次: 应报"与库里最新那份完全相同, 未新增" +make changes # 页面同款结论 +``` + +第一份快照落下之前, 变化提示一条都不会报 —— 这是设计如此 (`first`), 不是坏了。 +`make probe` (不加 `SNAP=1`) 是纯只读的, 随时可跑; 它**先比后写**, 所以加了 `SNAP=1` 也 +不会拿这份计划跟刚写进去的它自己比。 + +隔一天再跑一次 `make changes`, 就能看到第一份真正的跨日比对 —— 那时候才验得到 `cross_day` +与持仓票视图。账本还空着, 所以持仓那段现在必然是空的, 等账本重建后才有内容。 diff --git a/app/core/plan_diff.py b/app/core/plan_diff.py new file mode 100644 index 0000000..6c79406 --- /dev/null +++ b/app/core/plan_diff.py @@ -0,0 +1,334 @@ +# -*- coding: utf-8 -*- +""" +上游榜单的版本比对 (纯逻辑, 零外部依赖) +======================================== +上游 `/plan` 应答里有个 `changes` 字段, 本该回答「跟上一份比, 谁新进榜、谁掉榜」。 +实测它**恒为 `null`** (UPSTREAM_PLAN_API.md §4 Q4: 已问, 未给非空示例)。等它补, 就是把 +一件 PMS 自己能算的事挂在外部依赖上 —— 于是改成 PMS 每次拉计划落一份名册快照, 差异自己算。 + +顺带补上 §7.3 那个洞: `/plan` 是 `plan.collect()` **实时算**的, 同一个 `date` 可以对应 +多个版本 (07-31 实测: 相隔一小时的两次请求, 打分池 972/81 → 423/502, 榜首换人), 而应答里 +没有 `generated_at` 或版本戳 —— **下游此前没有任何办法判断"手上这份是哪一版"**。落了快照, +这个问题就地解决: 同一 `plan_date` 底下有几行, 就是上游那天重算过几版。 + +三种比对语义, 混在一起看会得出错误结论, 所以显式分开 (`kind`): + + * `first` —— 没有上一份快照。**一条变化都不报** (否则首次全表 300 条"新进榜") + * `cross_day` —— 计划日不同。这是 Q4 说的「每日几十只」的正常升降档 + * `same_date_revision` —— **同一个计划日的不同版本**。这不是市场变化, 是上游重算; + 它意味着 PMS 一小时前拿到的候选池与现在不是一回事, 要显眼报出来 + +-------------------------------------------------------------------------------- +掉榜为什么不能直接算「上一版有、这一版没有」 +-------------------------------------------------------------------------------- +我们向上游要 `top=300`。一只票从 rank 250 掉到 rank 310, 它会从名册里"消失"—— 但这不是 +上游观点变化, 只是排到了我们要的条数之外。**吃满 top 的那一版, 尾部的进出全是截断噪音。** +不作区分的话, 每天会多出几十上百条假掉榜, 提示就没人看了 (同 `_capped` 治的那个假警报, +是同一类错误的两个面)。 + +新进榜是**对称**的同一个问题: 上一版吃满了 top, 那么这一版排在末尾的"新面孔"很可能上一版 +就在榜上, 只是排在了 300 名开外看不见。所以两道闸各看各的那一版: + + 掉榜 —— 看 **这一版** 吃没吃满 (它截断了, 掉出去的可能只是被截掉) + 新进 —— 看 **上一版** 吃没吃满 (它截断了, 新面孔可能上一版就在, 只是没露面) + +`tail_guard` 是那道闸的位置: 排名进入榜单前 `tail_guard` 比例 (默认一半) 才算数, 其余计进 +`tail_churn` 只报个数、不列名。没吃满 top 的那一版不设闸 —— 那种情况下的进出是真的。 + +**持仓票不受这道闸约束。** 闸是拿"少报几条噪音"换"提示还有人看", 这笔账在持仓票上不成立: +非持仓票误报一条, 代价是白看一眼; **持仓票漏报一条, 代价是一个该减没减的仓位**。所以持仓票 +一律照报, 并在 `tail_churn.held_exempt` 里记下有几条是靠这条豁免出来的。 +""" +from __future__ import annotations + +import hashlib + +# 传导档位的强弱次序 (上游 2026-07-31 答复 Q7: 强传导 > 弱传导 > 无传导)。 +# 不在表里的档位 **不参与升降判定** —— 上游哪天加了新档位, 宁可不报, 不能报反。 +TIER_ORDER = {"强传导": 3, "弱传导": 2, "无传导": 1} + +BUCKET_MAIN, BUCKET_OBSERVE = "main", "observe" + +KIND_FIRST, KIND_CROSS_DAY, KIND_REVISION = "first", "cross_day", "same_date_revision" + +DEFAULT_RANK_JUMP = 50 # 名次跳变多少名才值得报 +DEFAULT_TAIL_GUARD = 0.5 # 吃满 top 时, 上一版排名在前多少比例内的掉榜才算数 + +# 名册里每条存哪些字段。键名取短的 —— 一份名册 300~400 条, 要落进 MEDIUMTEXT。 +_FIELDS = ("c", "n", "r", "s", "t", "h", "b") + + +def tier_rank(tier): + """档位 → 强弱序数; 未知档位返回 None (判不了升降就不判)。""" + if tier is None: + return None + return TIER_ORDER.get(str(tier).strip()) + + +# ================================================================ 名册 +def roster_of(plan: dict) -> dict: + """`plan_feed.parse_plan()` 的结果 → 名册 {ts_code: 精简行}。 + + 主榜在前, 同一只票若两档都出现以主榜为准 (与 parse_plan 里 themes 的口径一致)。 + """ + out = {} + for bucket in (BUCKET_MAIN, BUCKET_OBSERVE): + for r in (plan.get(bucket) or []): + code = r.get("ts_code") + if not code or code in out: + continue + out[code] = {"c": code, "n": r.get("name"), "r": r.get("rank"), + "s": r.get("score"), "t": r.get("tier"), + "h": r.get("theme"), "b": r.get("bucket") or bucket} + return out + + +def roster_rows(roster: dict) -> list: + """名册 → 落库用的 list (代码序, 保证同一份名册序列化结果稳定)。""" + return [{k: roster[c].get(k) for k in _FIELDS} for c in sorted(roster)] + + +def roster_from_rows(rows) -> dict: + """落库 list → 名册 dict。坏行跳过, 不让一条脏数据废掉整份快照。""" + out = {} + for r in (rows or []): + if not isinstance(r, dict): + continue + code = r.get("c") or r.get("ts_code") + if not code: + continue + out[code] = {"c": code, "n": r.get("n"), "r": r.get("r"), "s": r.get("s"), + "t": r.get("t"), "h": r.get("h"), "b": r.get("b")} + return out + + +def digest_of(roster: dict) -> str: + """名册指纹 —— 同一份榜重复拉不重复落。 + + **不含 score。** `/plan` 是实时算的, score 末位天天抖; 把它算进指纹, 每次缓存过期重拉 + 都会多落一行快照, 表白涨而信息量为零。含 rank 就够了: 分数抖到改变了次序才算真的变了, + 没改变次序的抖动对候选池没有任何影响 (候选就是按 score 降序切前 N)。 + """ + h = hashlib.sha1() + for c in sorted(roster): + r = roster[c] + h.update(("%s|%s|%s|%s|%s\n" % (c, r.get("r"), r.get("t") or "", + r.get("h") or "", r.get("b") or "")).encode("utf-8")) + return h.hexdigest() + + +# ================================================================ 比对 +def _row(r: dict, **extra) -> dict: + out = {"ts_code": r.get("c"), "name": r.get("n"), "rank": r.get("r"), + "score": r.get("s"), "tier": r.get("t"), "theme": r.get("h"), + "bucket": r.get("b")} + out.update(extra) + return out + + +def _bucket_len(roster: dict, bucket: str) -> int: + return sum(1 for r in (roster or {}).values() if (r.get("b") or BUCKET_MAIN) == bucket) + + +def _cutoff(roster: dict, capped, ratio: float) -> dict: + """尾部闸的位置 {bucket: 名次上限 or None}。 + + 只有那一版**吃满了 top** 才设闸 —— 没吃满说明上游能给的都给了, 榜尾的进出是真的。 + `capped` 按档给 (`{"main": bool, "observe": bool}`): 两档是两个独立的请求参数 + (`top` / `obs_top`), 主榜吃满**完全不意味着**观察档也吃满了。早先只传主榜那一个标志, + 效果是拿主榜的截断去解释观察档的进出 —— 观察档一次真实的摘牌会被记成"截断噪音"。 + 该档一条都没有时也不设闸: 整档消失/整档出现本身就是要报的大事, 不该被闸吞掉。 + """ + caps = capped if isinstance(capped, dict) else {BUCKET_MAIN: bool(capped), + BUCKET_OBSERVE: bool(capped)} + out = {} + for b in (BUCKET_MAIN, BUCKET_OBSERVE): + n = _bucket_len(roster, b) + out[b] = (n * ratio) if (caps.get(b) and n > 0) else None + return out + + +def _tail_dropped(row: dict, rank, cut: dict) -> bool: + """这条进/出是不是只是榜尾截断的噪音。 + + rank 缺失时返回 False = **照报**。全模块的口径是"判不了就不判"—— 在别处那意味着不报 + 变化, 在这里却要反过来: 这个函数的"是"代表**吞掉一条变化**, 所以拿不准时必须放行。 + """ + lim = cut.get(row.get("b") or BUCKET_MAIN) + return lim is not None and isinstance(rank, int) and rank > lim + + +def diff(prev_roster, curr_roster, *, prev_date=None, curr_date=None, held=(), + rank_jump: int = DEFAULT_RANK_JUMP, tail_guard: float = DEFAULT_TAIL_GUARD, + prev_capped=False, curr_capped=False, prev_broken: bool = False) -> dict: + """两份名册的差异。**不抛错** —— 页面提示不该有能力搞崩取数。 + + held: 当前持仓代码 (点式, 调用方负责归一)。持仓票的变化单独拎出来, 且**不受尾部闸约束** + —— 上游把一只持仓票摘出榜或降了档, 是"该不该继续拿着"的直接信号, 漏报的代价比误报大 + 得多 (见模块头部)。 + + prev_capped / curr_capped: 那一版吃没吃满请求的条数。可以给 bool, 也可以按档给 + `{"main": ..., "observe": ...}` —— 两档是两个独立的请求参数, 主榜吃满不代表观察档也满。 + 掉榜看 `curr_capped`、新进看 `prev_capped`, 各看各的那一版, 理由见模块头部。 + + prev_broken: 上一版**存在但读不出来** (名册落库时坏了)。这跟"没有上一版"是两回事: + 后者是正常的首次, 前者是故障 —— 必须说成故障, 否则一条真实变化都不报还显示得很正常。 + """ + held = {c for c in (held or []) if c} + curr_roster = curr_roster or {} + if not prev_roster: + note = ("上一版名册读不出来 (快照损坏?) —— **这次比不了**, 下面的空白不代表没有变化" + if prev_broken else + "首次落快照, 无可比对的上一版 —— 变化提示从下一次拉取开始") + return _empty(KIND_FIRST, prev_date, curr_date, curr_roster, held, note=note, + rank_jump=rank_jump, tail_guard=tail_guard, prev_broken=prev_broken) + + kind = (KIND_REVISION if (prev_date and curr_date and prev_date == curr_date) + else KIND_CROSS_DAY) + + entered, exited = [], [] + tail_churn = {"entered": 0, "exited": 0, "held_exempt": 0} + tier_up, tier_down, bucket_moved, jumps = [], [], [], [] + + # 两道尾部闸各看各的那一版 —— 见模块头部 + ratio = max(0.0, min(1.0, float(tail_guard if tail_guard is not None else DEFAULT_TAIL_GUARD))) + exit_cut = _cutoff(curr_roster, curr_capped, ratio) # 掉榜: 这一版截没截 + enter_cut = _cutoff(prev_roster, prev_capped, ratio) # 新进: 上一版截没截 + + for code, p in (prev_roster or {}).items(): + if curr_roster.get(code) is not None: + continue + h = code in held + if _tail_dropped(p, p.get("r"), exit_cut): + if not h: + tail_churn["exited"] += 1 # 本来就在榜尾 —— 大概率只是被这一版的 top 截掉 + continue + tail_churn["held_exempt"] += 1 # 持仓票不适用这道闸: 漏报一条就是一个仓位 + exited.append(_row(p, held=h)) + + for code, c in curr_roster.items(): + p = (prev_roster or {}).get(code) + if p is None: + h = code in held + if _tail_dropped(c, c.get("r"), enter_cut): + if not h: + tail_churn["entered"] += 1 # 上一版就吃满了, 这张脸当时可能只是没露面 + continue + tail_churn["held_exempt"] += 1 + entered.append(_row(c, held=h)) + continue + h = code in held + pb, cb = p.get("b") or BUCKET_MAIN, c.get("b") or BUCKET_MAIN + if pb != cb: + bucket_moved.append(_row(c, held=h, moved_from=pb, moved_to=cb)) + pt, ct = tier_rank(p.get("t")), tier_rank(c.get("t")) + if pt is not None and ct is not None and pt != ct: + (tier_up if ct > pt else tier_down).append( + _row(c, held=h, tier_from=p.get("t"), tier_to=c.get("t"))) + pr, cr = p.get("r"), c.get("r") + if pb == cb and isinstance(pr, int) and isinstance(cr, int): + d = cr - pr + if abs(d) >= max(1, int(rank_jump or DEFAULT_RANK_JUMP)): + jumps.append(_row(c, held=h, rank_from=pr, rank_to=cr, rank_delta=d)) + + entered.sort(key=lambda x: (x["rank"] if isinstance(x["rank"], int) else 10 ** 9)) + exited.sort(key=lambda x: (x["rank"] if isinstance(x["rank"], int) else 10 ** 9)) + jumps.sort(key=lambda x: abs(x.get("rank_delta") or 0), reverse=True) + for lst in (tier_up, tier_down, bucket_moved): + lst.sort(key=lambda x: (x["rank"] if isinstance(x["rank"], int) else 10 ** 9)) + + out = {"kind": kind, "prev_date": prev_date, "curr_date": curr_date, + "entered": entered, "exited": exited, "tier_up": tier_up, "tier_down": tier_down, + "bucket_moved": bucket_moved, "rank_jump": jumps, "tail_churn": tail_churn, + "tail_guarded": {"entered": _guard_view(enter_cut), "exited": _guard_view(exit_cut)}, + "roster_size": {"prev": len(prev_roster or {}), "curr": len(curr_roster)}, + "prev_broken": False, "date_regressed": _regressed(prev_date, curr_date), + "params": {"rank_jump": int(rank_jump or DEFAULT_RANK_JUMP), "tail_guard": ratio}} + out["counts"] = {k: len(out[k]) for k in + ("entered", "exited", "tier_up", "tier_down", "bucket_moved", "rank_jump")} + out["held"] = _held_view(out, held) + out["note"] = _note(out) + return out + + +def _guard_view(cut: dict) -> dict: + """哪几档真的设了闸, 闸位在第几名 —— 让"为什么这条没报"可查, 而不是只能猜。""" + return {b: (None if cut.get(b) is None else round(cut[b], 1)) + for b in (BUCKET_MAIN, BUCKET_OBSERVE)} + + +def _regressed(prev_date, curr_date) -> bool: + """这一版的计划日比上一版还早。 + + 正常流程不会这样 —— 快照按落库次序排, 而计划日是往前走的。会出现只有一种情况: + 有人拿 `--date` 显式拉了一份旧计划。那时"新进/掉榜"的方向是反的, 得说明白。 + """ + return bool(prev_date and curr_date and str(curr_date) < str(prev_date)) + + +def _empty(kind, prev_date, curr_date, curr_roster, held, note="", + rank_jump=DEFAULT_RANK_JUMP, tail_guard=DEFAULT_TAIL_GUARD, + prev_broken=False) -> dict: + ratio = max(0.0, min(1.0, float(tail_guard if tail_guard is not None else DEFAULT_TAIL_GUARD))) + out = {"kind": kind, "prev_date": prev_date, "curr_date": curr_date, + "entered": [], "exited": [], "tier_up": [], "tier_down": [], + "bucket_moved": [], "rank_jump": [], + "tail_churn": {"entered": 0, "exited": 0, "held_exempt": 0}, + "tail_guarded": {"entered": _guard_view({}), "exited": _guard_view({})}, + "roster_size": {"prev": 0, "curr": len(curr_roster or {})}, + "prev_broken": bool(prev_broken), "date_regressed": _regressed(prev_date, curr_date), + # 回显**调用方实际给的**参数, 不是默认值 —— 页面拿它显示当前口径, 写死会骗人 + "params": {"rank_jump": int(rank_jump or DEFAULT_RANK_JUMP), "tail_guard": ratio}} + out["counts"] = {k: 0 for k in + ("entered", "exited", "tier_up", "tier_down", "bucket_moved", "rank_jump")} + out["held"] = _held_view(out, held) + out["note"] = note + return out + + +def _held_view(d: dict, held) -> dict: + """持仓票命中的变化。**掉榜与降档排在最前** —— 那是要不要继续持有的信号。 + + `n_watch` 是"要盯的"那几条, 页面红条与探活脚本都拿它当判据。三类算进去: + + 掉榜 · 降档 · **主榜→观察档** + + 第三类容易漏。观察档的定义是"无券商覆盖", 所以掉进观察档的含义是**估值锚没了**, 与 + 降档是同一量级的信号; 而观察档的行没有 `tier`, 于是它既不算 `exited` 也不算 + `tier_down` —— 不显式算进来的话, 一只持仓票丢了券商覆盖, 页面红条不亮、探活脚本还会 + 打一句"持仓票没被摘也没降档"。 + """ + pick = lambda k: [x for x in d.get(k) or [] if x.get("held")] # noqa: E731 + out = {k: pick(k) for k in ("exited", "tier_down", "bucket_moved", + "tier_up", "entered", "rank_jump")} + out["coverage_lost"] = [x for x in out["bucket_moved"] if x.get("moved_to") == BUCKET_OBSERVE] + out["n_watch"] = len(out["exited"]) + len(out["tier_down"]) + len(out["coverage_lost"]) + out["n_total"] = sum(len(v) for k, v in out.items() + if isinstance(v, list) and k != "coverage_lost") # 别重复计 bucket_moved + return out + + +def _note(d: dict) -> str: + """一句人话结论 —— 页面横幅与日志直接用这句, 不必各写一份措辞。""" + c, h = d["counts"], d["held"] + if d["kind"] == KIND_REVISION: + head = (f"同一计划日 {d['curr_date']} 的**新版本**: 上游重算过 —— " + f"此前拿到的候选池与现在不是同一份") + else: + head = f"计划日 {d['prev_date']} → {d['curr_date']}" + body = (f"新进 {c['entered']} · 掉榜 {c['exited']} · 升档 {c['tier_up']} · " + f"降档 {c['tier_down']} · 覆盖翻转 {c['bucket_moved']}") + tail = "" + if h["n_watch"]: + tail = f" · **持仓票 {h['n_watch']} 只掉榜/降档/丢了券商覆盖**" + elif h["n_total"]: + tail = f" · 持仓票命中 {h['n_total']} 条" + churn = d.get("tail_churn") or {} + n_churn = int(churn.get("entered") or 0) + int(churn.get("exited") or 0) + if n_churn: + tail += f" · 榜尾进出 {n_churn} 只未计入 (那一版吃满 top, 尾部不可信)" + if churn.get("held_exempt"): + tail += f" · 其中 {churn['held_exempt']} 只因是持仓票照报" + if d.get("date_regressed"): + tail += " · **注意计划日是往回走的** (拉了一份旧计划?), 新进/掉榜的方向是反的" + return f"{head}: {body}{tail}" diff --git a/app/repo/pms_repo.py b/app/repo/pms_repo.py index 0002ff4..bd8da71 100644 --- a/app/repo/pms_repo.py +++ b/app/repo/pms_repo.py @@ -573,3 +573,99 @@ def upsert_industry(rows: list) -> int: return execute_many( "INSERT INTO pms_industry_map (ts_code, industry, updated_at) VALUES (:code, :ind, :ts) " "ON DUPLICATE KEY UPDATE industry = :ind, updated_at = :ts", payload) + + +# ================================================================ pms_plan_snapshot +def latest_plan_digest(): + """最新一份快照的 (plan_date, digest); 没有则 None。判"这份榜变没变"只该拿它比。""" + r = fetch_one("SELECT plan_date, digest FROM pms_plan_snapshot ORDER BY id DESC LIMIT 1") + return (r["plan_date"], r["digest"]) if r else None + + +def insert_plan_snapshot(*, plan_date, digest, roster, meta=None, n_main=0, n_observe=0, + capped_main=False, capped_obs=False, fetched_at=None) -> int: + """落一份名册快照; 内容与**最新那份**相同则跳过, 返回 0。 + + 两个坑都在这一句"跟最新那份比"上, 值得写清楚: + + 1. **不能拿 rowcount 判是不是新插的。** 本文件 `inbox_put` 那里已经栽过一次: + SQLAlchemy 的 MySQL 方言默认开 `CLIENT_FOUND_ROWS`, `ON DUPLICATE KEY UPDATE id = id` + 命中重复照样回 1。所以先查后插, 由**查的结果**决定 `stored`, 不看 rowcount。 + 2. **去重的对象是"最新那份", 不是"曾经出现过的任何一份"。** 早先给 + `(plan_date, digest)` 加了唯一键当去重手段, 结果是: 同一天上游 A→B→A 地改回去时, + 第三次的 A 因为"曾经存过"而落不下去, 库里最新仍是 B —— 于是页面把 B→A 这次**真实的 + 回退**显示成 A→B, 方向整个反了, 还会顺带报一句"手上这份没落进快照 (落库失败?)"。 + 现在唯一键去掉, 只跟最新那行比: 变了就存, 没变就跳过。表因此可能出现同 digest 的 + 多行, 那正是想记的事实 —— 上游确实改回来过。 + + 并发下两个协程可能同时判定"变了"而落两行相同的快照。代价是多一行 (比对结果为"无变化", + 不产生任何假信号), 换掉的是一整类方向错误, 这笔交换是划算的。 + """ + now = fetched_at or _NOW() + cur = latest_plan_digest() + if cur and cur[0] == str(plan_date) and cur[1] == str(digest): + return 0 + execute( + "INSERT INTO pms_plan_snapshot " + " (plan_date, digest, fetched_ymd, fetched_at, n_main, n_observe, " + " capped_main, capped_obs, roster_json, meta_json) " + "VALUES (:d, :g, :ymd, :ts, :nm, :no, :cm, :co, :roster, :meta)", + {"d": str(plan_date), "g": str(digest), "ymd": int(now.strftime("%Y%m%d")), "ts": now, + "nm": int(n_main or 0), "no": int(n_observe or 0), + "cm": 1 if capped_main else 0, "co": 1 if capped_obs else 0, + "roster": _dumps(roster or []), "meta": _dumps(meta or {})}) + return 1 + + +def latest_plan_snapshots(limit: int = 2) -> list: + """最近若干份快照, **新的在前**。带名册 —— 比对要用。""" + rows = fetch_all( + "SELECT id, plan_date, digest, fetched_ymd, fetched_at, n_main, n_observe, " + " capped_main, capped_obs, roster_json, meta_json " + "FROM pms_plan_snapshot ORDER BY id DESC LIMIT :n", {"n": max(1, int(limit))}) + return [_snapshot_row(r) for r in rows] + + +def list_plan_snapshots(*, plan_date=None, limit: int = 50) -> list: + """快照台账 (不带名册 —— 页面列表不需要那 20KB)。 + + 同一 plan_date 有几行, 就是上游那天重算过几版 (UPSTREAM_PLAN_API.md §7.3)。 + """ + sql = ("SELECT id, plan_date, digest, fetched_ymd, fetched_at, n_main, n_observe, " + " capped_main, capped_obs, meta_json FROM pms_plan_snapshot ") + params = {"n": max(1, int(limit))} + if plan_date: + sql += "WHERE plan_date = :d " + params["d"] = str(plan_date) + rows = fetch_all(sql + "ORDER BY id DESC LIMIT :n", params) + for r in rows: + r["meta"] = _loads(r.pop("meta_json", None), {}) + r["capped_main"] = bool(r.get("capped_main")) # 与 latest_plan_snapshots 保持同型 + r["capped_obs"] = bool(r.get("capped_obs")) + return rows + + +def prune_plan_snapshots(keep: int = 200) -> int: + """只留最近 keep 份。 + + 单表守卫不许子查询跨表, 所以先查出保留水位的 id 再按 id 删 —— 两条单表语句, + 比一条 `DELETE ... WHERE id NOT IN (SELECT ...)` 干净, 也不会被守卫拦下。 + + 下限是 **2 份不是 1 份**: 比对天生要两份, 留一份等于把榜单变化永久锁死在"首次"上, + 而且页面上看起来一切正常。参数填了 0 或 1 也按 2 算。 + """ + keep = max(2, int(keep or 0)) + r = fetch_one("SELECT id FROM pms_plan_snapshot ORDER BY id DESC LIMIT 1 OFFSET :k", + {"k": keep}) + if not r: + return 0 + return execute("DELETE FROM pms_plan_snapshot WHERE id <= :id", {"id": int(r["id"])}) + + +def _snapshot_row(r: dict) -> dict: + out = dict(r) + out["roster"] = _loads(out.pop("roster_json", None), []) + out["meta"] = _loads(out.pop("meta_json", None), {}) + out["capped_main"] = bool(out.get("capped_main")) + out["capped_obs"] = bool(out.get("capped_obs")) + return out diff --git a/app/services/plan_feed.py b/app/services/plan_feed.py index bfa8bb1..136b544 100644 --- a/app/services/plan_feed.py +++ b/app/services/plan_feed.py @@ -37,8 +37,15 @@ 所以 `PMS_PLAN_THEME_SYNC` 默认关, 行业源走 `gp_stock_category`; theme 只用在候选阶段的 `PMS_PLAN_THEME_CAP_LOCAL` 上 —— 防单一传导主题刷屏, 那才是它擅长的事。 +**榜单变化 PMS 自己算** (2026-07-31 加, 见 UPSTREAM_PLAN_API.md §9): 应答里的 `changes` +字段恒为 `null` (§4 Q4 已问未答), 而「谁新进榜、谁掉榜、谁降了档」是上游观点变化最直接的 +信号。与其等对方补, 不如每次拉计划落一份名册快照 (`pms_plan_snapshot`), 差异由 +`app.core.plan_diff` 纯逻辑算。顺带补上 §7.3 那个洞: `/plan` 是实时算的, 同一个 `date` 能 +对应多个版本而应答里没有版本戳 —— 落了快照, 同一 `plan_date` 底下有几行就是重算过几版。 + 模块级只依赖 stdlib + `app.core.command_spec` (纯逻辑), 其余 (requests / param_store / -pms_repo / tradedays) 一律函数内懒加载 —— 让解析与筛选这两段纯逻辑可以零依赖单测。 +pms_repo / plan_diff / tradedays) 一律函数内懒加载 —— 让解析与筛选这两段纯逻辑可以零依赖 +单测。 """ from __future__ import annotations @@ -342,6 +349,10 @@ def _params() -> dict: "exclude_st": ps.get_bool("PMS_PLAN_EXCLUDE_ST", True), "query_extra": parse_query_extra(ps.get("PMS_PLAN_QUERY_EXTRA", "")), "source": (ps.get("PMS_CANDIDATE_SOURCE", SRC_PLAN_API) or SRC_PLAN_API).strip(), + "snapshot": ps.get_bool("PMS_PLAN_SNAPSHOT", True), + "snapshot_keep": ps.get_int("PMS_PLAN_SNAPSHOT_KEEP", 200), + "diff_rank_jump": ps.get_int("PMS_PLAN_DIFF_RANK_JUMP", 50), + "diff_tail_guard": ps.get_float("PMS_PLAN_DIFF_TAIL_GUARD", 0.5), } @@ -451,6 +462,8 @@ def get_plan(*, force: bool = False, date=None) -> dict: plan["age_tdays"] = plan_age_tdays(plan["date"]) if p["theme_sync"]: plan["theme_sync"] = _sync_themes_quiet(plan) + if p["snapshot"]: + plan["snapshot"] = _snapshot_quiet(plan, keep=p["snapshot_keep"]) with _lock: _cache.update({"at": time.time(), "key": key, "plan": plan, "error": None}) logger.info("[上游计划] %s 主榜 %d / 观察 %d (日龄 %d 交易日) ← %s", @@ -490,6 +503,207 @@ def _sync_themes_quiet(plan: dict) -> dict: return {"rows": 0, "affected": 0, "error": f"{type(e).__name__}: {e}"} +# ================================================================ 名册快照与榜单变化 +def snapshot(plan: dict, *, keep: int = 200) -> dict: + """把这一份计划的名册落进 pms_plan_snapshot。 + + 同一份榜重复拉不重复落 (靠 uk_date_digest)。返回里带 `stored` —— 落进去了才是新版本, + 没落进去说明这份榜跟库里最新那份一模一样。 + """ + from app.core import plan_diff as pdf + from app.repo import pms_repo + roster = pdf.roster_of(plan) + digest = pdf.digest_of(roster) + meta = {k: plan.get(k) for k in ("counts", "returned", "requested", "funnel", + "theme_cap", "heat_date", "market_snapshot_days", + "age_tdays", "url")} + truncated = plan.get("truncated") or {} + n = pms_repo.insert_plan_snapshot( + plan_date=plan.get("date"), digest=digest, roster=pdf.roster_rows(roster), meta=meta, + n_main=(plan.get("returned") or {}).get("main") or 0, + n_observe=(plan.get("returned") or {}).get("observe") or 0, + capped_main=bool(truncated.get("main")), capped_obs=bool(truncated.get("observe"))) + out = {"digest": digest, "rows": len(roster), "stored": bool(n)} + if n: + out["pruned"] = pms_repo.prune_plan_snapshots(keep=keep) + logger.info("[上游计划] 名册快照已落: %s digest=%s (%d 只)", + plan.get("date"), digest[:10], len(roster)) + return out + + +def _snapshot_quiet(plan: dict, *, keep: int = 200) -> dict: + """快照落库失败**不许阻断候选池** —— 没有变化提示是可接受的降级, 没候选不是。 + + 同 `_sync_themes_quiet` 的口径: 这是个"锦上添花"的旁路, 它的故障不该传染主链路。 + 但也不静默: error 会一路带到页面上, 因为快照断了就意味着变化提示从此刻起是错的 + (拿旧版本当上一版比)。 + """ + try: + return snapshot(plan, keep=keep) + except Exception as e: + logger.warning("[上游计划] 名册快照落库失败 (榜单变化提示将不可用): %s", e) + return {"stored": False, "error": f"{type(e).__name__}: {e}"} + + +def _held_codes() -> list: + """当前持仓代码。拿不到就当空 —— 持仓视图缺失只是少一段提示, 不该让整个接口失败。""" + try: + from app.repo import pms_repo + return [p["ts_code"] for p in pms_repo.list_positions(only_open=True)] + except Exception as e: + logger.debug("[上游计划] 取持仓失败, 榜单变化的持仓视图留空: %s", e) + return [] + + +def _capped_of(row) -> dict: + """一行快照的按档截断标志。两档是两个独立的请求参数 (`top` / `obs_top`), 各判各的 —— + 拿主榜的截断去解释观察档的进出, 会把观察档一次真实的摘牌记成"截断噪音"。""" + row = row or {} + return {"main": bool(row.get("capped_main")), "observe": bool(row.get("capped_obs"))} + + +def _prev_broken(prev) -> bool: + """上一行"说自己有几百只、名册却解析成空" = 落库时坏了, 不是正常的"没有上一版"。 + + `pms_repo._loads` 在 JSON 解析失败时返回默认值 `[]`, 于是坏数据长得跟"首次"一模一样 + —— 一条真实变化都不报, 页面还显示得很正常。这里把两者分开。 + """ + if not prev: + return False + claimed = int(prev.get("n_main") or 0) + int(prev.get("n_observe") or 0) + return bool(claimed > 0 and not (prev.get("roster") or [])) + + +def _diff_against(prev, curr_roster, *, curr_date, curr_capped, held, rank_jump, tail_guard): + """比对的公共落点 —— `changes` (库里两份) 与 `preview_changes` (手上这份 vs 库里最新) + 只在"谁是 curr"上不同, 比对口径必须完全一致, 所以共用这一段。""" + from app.core import command_spec as cs + from app.core import plan_diff as pdf + p = _params() + codes = held if held is not None else _held_codes() + # 名册的键是点式 (parse_plan 已归一)。持仓来源万一给的是前缀式, 不归一的后果是 + # **持仓视图静默为空** —— 不报错、不少数据, 就是一条都不命中。 + codes = [cs.normalize_code(str(c)) for c in (codes or []) if c] + d = pdf.diff(pdf.roster_from_rows(prev["roster"]) if prev else None, curr_roster, + prev_date=(prev or {}).get("plan_date"), curr_date=curr_date, + held=codes, + rank_jump=(rank_jump if rank_jump is not None else p["diff_rank_jump"]), + tail_guard=(tail_guard if tail_guard is not None else p["diff_tail_guard"]), + prev_capped=_capped_of(prev), curr_capped=curr_capped, + prev_broken=_prev_broken(prev)) + d["prev_fetched_at"] = (prev or {}).get("fetched_at") + return d + + +def preview_changes(plan: dict, *, held=None, rank_jump=None, tail_guard=None) -> dict: + """**只读**: 拿手上这份计划跟库里最新那份快照比, 一个字都不写。 + + 探活脚本用 —— 那个脚本的契约是"只读, 不写任何表", 破了它比少个功能糟得多。 + 页面走 `changes()`: 页面本来就要落快照 (那是它的正常职责), 比的是库里最近两份。 + """ + from app.core import plan_diff as pdf + out = {"ok": False, "readonly": True, "hint": "", "diff": None, "stored": None} + try: + from app.repo import pms_repo + snaps = pms_repo.latest_plan_snapshots(limit=1) + except Exception as e: + out["hint"] = f"读名册快照失败: {type(e).__name__}: {e}" + return out + prev = snaps[0] if snaps else None + curr_roster = pdf.roster_of(plan) + out["stored"] = ({k: prev.get(k) for k in ("id", "plan_date", "digest", "fetched_at")} + if prev else None) + out["same_as_stored"] = bool(prev and prev["digest"] == pdf.digest_of(curr_roster)) + try: + tr = plan.get("truncated") or {} + out["diff"] = _diff_against(prev, curr_roster, curr_date=plan.get("date"), + curr_capped={"main": bool(tr.get("main")), + "observe": bool(tr.get("observe"))}, + held=held, rank_jump=rank_jump, tail_guard=tail_guard) + except Exception as e: + out["hint"] = f"比对失败: {type(e).__name__}: {e}" + return out + out["ok"] = True + out["hint"] = out["diff"]["note"] + return out + + +def changes(*, held=None, rank_jump=None, tail_guard=None) -> dict: + """榜单变化: 拿库里最近两份快照比。**绝不抛错** —— 这是页面提示, 不该有能力搞崩取数。 + + 上游 `changes` 字段恒为 null (UPSTREAM_PLAN_API.md §4 Q4), 所以这段由 PMS 自己算。 + 比对语义与尾部闸见 `app/core/plan_diff.py` 的模块说明。 + + `in_sync`: 手上这份计划的指纹与库里最新快照**是不是同一个**。落库失败过 (DB 抖了一下) + 的话这里会是 False —— 那说明比出来的是两个旧版本, 跟你现在用的候选池不是一回事。 + 宁可把这句话摆出来, 也不让页面显示一个看起来正常、实际过时的差异。 + """ + from app.core import plan_diff as pdf + p = _params() + out = {"ok": False, "enabled": bool(p["snapshot"]), "hint": "", "diff": None, + "snapshots": [], "in_sync": None} + if not p["snapshot"]: + out["hint"] = "名册快照已关 (PMS_PLAN_SNAPSHOT=false) —— 榜单变化提示不可用" + return out + + curr_digest = None + try: # 先确保手上这份已落库 (缓存命中时是个空操作) + plan = get_plan() + curr_digest = pdf.digest_of(pdf.roster_of(plan)) + except PlanFeedError as e: + out["hint"] = f"上游计划取不到, 只能拿库里的历史快照比: {e}" + except Exception as e: + out["hint"] = f"{type(e).__name__}: {e}" + + try: + from app.repo import pms_repo + snaps = pms_repo.latest_plan_snapshots(limit=2) + except Exception as e: + out["hint"] = f"读名册快照失败: {type(e).__name__}: {e}" + return out + + out["snapshots"] = [{k: s.get(k) for k in + ("id", "plan_date", "digest", "fetched_at", "n_main", "n_observe", + "capped_main", "capped_obs")} for s in snaps] + if not snaps: + out.update({"ok": True, "hint": out["hint"] or "还没有任何名册快照 —— 拉一次计划即有"}) + return out + if curr_digest is not None: + out["in_sync"] = (snaps[0]["digest"] == curr_digest) + + curr, prev = snaps[0], (snaps[1] if len(snaps) > 1 else None) + try: + d = _diff_against(prev, pdf.roster_from_rows(curr["roster"]), + curr_date=curr.get("plan_date"), curr_capped=_capped_of(curr), + held=held, rank_jump=rank_jump, tail_guard=tail_guard) + except Exception as e: # 兑现"绝不抛错": 顶层 ok 一为假页面就挂全局红条 + logger.warning("[上游计划] 榜单变化比对失败: %s", e) + out["hint"] = f"比对失败: {type(e).__name__}: {e}" + return out + d["curr_fetched_at"] = curr.get("fetched_at") + out.update({"ok": True, "diff": d}) + if out["in_sync"] is False: + out["hint"] = ("手上这份计划**没落进快照** (上次落库失败?) —— 下面比的是两个旧版本, " + "与当前候选池不是同一份。" + (out["hint"] or "")) + elif not out["hint"]: + out["hint"] = d["note"] + return out + + +def snapshot_log(*, plan_date=None, limit: int = 50) -> dict: + """快照台账。同一 `plan_date` 有几行, 就是上游那天重算过几版 (§7.3)。""" + try: + from app.repo import pms_repo + rows = pms_repo.list_plan_snapshots(plan_date=plan_date, limit=limit) + by_date = {} + for r in rows: + by_date[r["plan_date"]] = by_date.get(r["plan_date"], 0) + 1 + revised = {k: v for k, v in by_date.items() if v > 1} + return {"ok": True, "rows": rows, "versions_per_date": by_date, "revised": revised} + except Exception as e: + return {"ok": False, "rows": [], "error": f"{type(e).__name__}: {e}"} + + # ================================================================ 对外: 候选与状态 def candidates(*, held=(), black=()) -> dict: """参数驱动的候选清单 (不含价格 —— 价格由调用方用 market.get_price 现取)。""" @@ -513,7 +727,8 @@ def status() -> dict: "top_n": p["top_n"], "tiers": p["tiers"], "min_upside": p["min_upside"], "theme_cap_local": p["theme_cap_local"], "query": build_query(p), "include_observe": p["include_observe"], "stale_tdays": p["stale_tdays"], - "theme_sync": p["theme_sync"], "enabled": bool(p["base"])} + "theme_sync": p["theme_sync"], "snapshot_on": p["snapshot"], + "enabled": bool(p["base"])} if not p["base"]: st.update({"ok": False, "hint": "上游计划接口未配置 (PMS_PLAN_API_BASE 为空) —— " "候选池将为空, 升仓/建仓类命令无票可选"}) @@ -533,7 +748,7 @@ def status() -> dict: "funnel": plan.get("funnel"), "requested": plan.get("requested"), "truncated": plan.get("truncated"), "theme_cap": plan.get("theme_cap"), "encoding": plan.get("encoding"), - "theme_sync": plan.get("theme_sync"), + "theme_sync": plan.get("theme_sync"), "snapshot": plan.get("snapshot"), "fetched_at": plan.get("fetched_at"), "url": plan.get("url"), "hint": f"计划 {plan['date']} 已就绪 (日龄 {plan.get('age_tdays')} 交易日)"}) return st diff --git a/app/web/main.py b/app/web/main.py index 12d6bce..5918ca5 100644 --- a/app/web/main.py +++ b/app/web/main.py @@ -439,10 +439,27 @@ def api_plan_refresh(date: str = Query(None)): plan = plan_feed.get_plan(force=True, date=(date or None)) return {"ok": True, "date": plan["date"], "age_tdays": plan.get("age_tdays"), "returned": plan["returned"], "theme_sync": plan.get("theme_sync"), - "industry": industry.status()} + "snapshot": plan.get("snapshot"), "industry": industry.status()} return ok(_refresh) +@app.get("/api/upstream/plan-changes") +def api_plan_changes(limit: int = Query(50)): + """榜单变化 (PMS 自算, 上游 `changes` 字段恒为 null —— UPSTREAM_PLAN_API.md §9)。 + + 绝不抛错: 变化提示是锦上添花, 它坏了不该让「上游计划」抽屉打不开。 + `log.revised` 里某个 plan_date 出现多版, 就是 §7.3 那个"同一 date 会变"。 + + 子状态**嵌在 changes 里**而不是顶层 —— 顶层 `ok` 一为假, 页面会挂全局红条。 + 「快照关着」「还没有快照」都不是接口故障, 不该长成那副样子 (同 status 的处理)。 + """ + def _changes(): + from app.services import plan_feed + return {"ok": True, "changes": plan_feed.changes(), + "log": plan_feed.snapshot_log(limit=max(1, int(limit or 50)))} + return ok(_changes) + + @app.get("/api/ops/downstream-schema") def api_downstream_schema(): """导出下游三表的实际列定义 —— 用于回填 QMT_INTERFACE_REQUIREMENTS D1。""" diff --git a/app/web/static/index.html b/app/web/static/index.html index 6ef198e..9cab865 100644 --- a/app/web/static/index.html +++ b/app/web/static/index.html @@ -576,6 +576,136 @@ {{ planCand.st_unknown.length }} 只无名称, ST 判不了已放行
{{ (planCand.items||[]).map(x=>x.ts_code).join(' ') }} + + 榜单变化 +
+ 重算变化 + {{ chgSt.hint }} +
+ + + + + + + + + + +
+ {{ chg.kind === 'first' ? '首次落快照' : (chg.prev_date + ' → ' + chg.curr_date) }} · + 名册 {{ chg.roster_size.prev }} → {{ chg.roster_size.curr }} · + 新进 {{ chg.counts.entered }} · 掉榜 {{ chg.counts.exited }} · + 升档 {{ chg.counts.tier_up }} · 降档 {{ chg.counts.tier_down }} · + 覆盖翻转 {{ chg.counts.bucket_moved }} · 名次跳变 {{ chg.counts.rank_jump }} + · + 榜尾进出 {{ chg.tail_churn.entered + chg.tail_churn.exited }} 只未计入 (那一版吃满 top, + 尾部是截断噪音不是观点变化) + · 其中 + {{ chg.tail_churn.held_exempt }} 只因是持仓票照报 (持仓票不受榜尾闸约束 —— + 漏报一条就是一个仓位) + · 同一计划日出现多版: + {{ chgRevised.map(x => x[0] + '×' + x[1]).join(', ') }} +
+ + + + + + + + + + + +
+ 上游把一只持仓票摘出榜或降了档, 是"还该不该继续拿着"的直接信号 —— + 单独拎出来免得被几十条新进榜淹掉。它只是提示, 不自动产生任何指令: + 要减要清仍走命令台或提议确认。 +
+
+ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + +
+ 同一「计划日」出现多行 = 上游那天重算过多版(UPSTREAM_PLAN_API.md §7.3: + `/plan` 是实时算的, 应答里没有 generated_at, 下游本来没办法判断手上这份是哪一版)。 + 指纹不含 score —— 分数抖动但次序没变不算新版本, 否则表白涨而信息量为零。 +
+
+
+ 主榜明细 @@ -642,6 +772,24 @@ createApp({ const planStatus = computed(() => planRaw.value.status || {}); const planRows = computed(() => planRaw.value.rows || []); const planCand = computed(() => planRaw.value.candidates_raw || null); + // 榜单变化 (PMS 自算 —— 上游 changes 字段恒为 null)。另起一个接口, 不塞进 /api/upstream/plan: + // 它要读库, 而计划预览要读上游, 两者的故障域不同, 混在一起会互相拖垮。 + const chgRaw = ref({}), chgLoading = ref(false), chgTab = ref('held'); + const chgSt = computed(() => chgRaw.value.changes || {}); + const chg = computed(() => chgSt.value.diff || null); + const chgHeld = computed(() => (chg.value || {}).held || {}); + const chgLog = computed(() => chgRaw.value.log || {}); + const chgRevised = computed(() => Object.entries(chgLog.value.revised || {})); + // 掉榜/降档合成一张表 —— 「上游把这只票摘了或降了档」是同一个意思的两种说法 + const chgWatch = computed(() => { + const d = chg.value; if (!d) return []; + // 三类都算「上游把这只票看淡了」: 掉榜 / 降档 / 掉进观察档(=没券商覆盖了, 估值锚没了)。 + // 与 plan_diff 的 held.n_watch 口径一致 —— 两边分家的话红条会亮而表里没有对应行。 + return (d.exited || []).map(x => ({ ...x, why: '掉榜' })) + .concat((d.tier_down || []).map(x => ({ ...x, why: '降档 ' + x.tier_from + '→' + x.tier_to }))) + .concat((d.bucket_moved || []).filter(x => x.moved_to === 'observe') + .map(x => ({ ...x, why: '券商覆盖没了 (主榜→观察档)' }))); + }); const issuing = ref(false); const form = reactive({ cmd_type: '', params: {}, note: '' }); @@ -900,13 +1048,20 @@ createApp({ planRaw.value = (d && d.ok) ? d : { status: { ok: false, hint: (d||{}).error || '取数失败' } }; planLoading.value = false; } - async function openPlan() { planDrawer.value = true; await loadPlan(); } + async function loadChanges() { + chgLoading.value = true; + const d = await call('get', '/api/upstream/plan-changes?limit=50'); + chgRaw.value = (d && d.ok) ? d + : { changes: { ok: false, hint: (d || {}).error || '取数失败' }, log: {} }; + chgLoading.value = false; + } + async function openPlan() { planDrawer.value = true; await Promise.all([loadPlan(), loadChanges()]); } async function refreshPlan() { planLoading.value = true; const d = await call('post', '/api/ops/plan-refresh'); planResult.value = JSON.stringify(d, null, 2); planLoading.value = false; - await Promise.all([loadPlan(), loadOverview()]); + await Promise.all([loadPlan(), loadChanges(), loadOverview()]); } async function openReport() { const d = await call('get', '/api/report'); report.value = d.data || d || {}; @@ -923,7 +1078,9 @@ createApp({ loadAll, loadParams, saveParams, loadPlans, loadLots, onCmdChange, issue, cancelCmd, replan, decide, ops, loadSchema, openReport, cancelIns, planDrawer, planLoading, planStatus, planRows, planCand, planResult, - loadPlan, openPlan, refreshPlan }; + loadPlan, openPlan, refreshPlan, + chgRaw, chgSt, chgLoading, chgTab, chg, chgHeld, chgLog, chgRevised, chgWatch, + loadChanges }; } }).use(ElementPlus).mount('#app'); diff --git a/config/settings.py b/config/settings.py index 471d340..e0bc39d 100644 --- a/config/settings.py +++ b/config/settings.py @@ -107,6 +107,15 @@ class Settings(BaseSettings): PMS_PLAN_QUERY_EXTRA: str = "" # 附加查询串逃生口, 如 "foo=1"。上游哪天加了新参数 # 不用改代码即可透传; top/obs_top/theme_cap 已有显式参数, 同名键以显式参数为准 + # --- 榜单变化提示 (PMS 自算; 上游 changes 字段恒为 null, 见 UPSTREAM_PLAN_API.md §9) --- + PMS_PLAN_SNAPSHOT: bool = True # 每次拉计划落一份名册快照 (变化提示的底本)。 + # 关掉只是没有变化提示, 不影响候选池; 但同时也失去"手上这份是哪一版"的自证能力 + PMS_PLAN_SNAPSHOT_KEEP: int = 200 # 快照保留份数 (一份 300 只约 20KB) + PMS_PLAN_DIFF_RANK_JUMP: int = 50 # 名次跳变多少名才报 + PMS_PLAN_DIFF_TAIL_GUARD: float = 0.5 # 榜尾进出的噪音闸: 排名进入榜单前这个比例才算数。 + # 榜是按 top 截断过的 —— 从 rank 250 掉到 310 只是被截掉, 不是上游摘了它。不设这道闸 + # 每天会多几十上百条假掉榜, 提示就没人看了。1.0=只挡榜长之外, 0=一律当噪音 + # --- 行业约束 (硬拦截; 数据源接口化) --- # 行业划分数据源。gp_hybk = 199 库的行业板块表 (三级 884*, 每票只认一个主行业), # 2026-07-31 起的正式口径; gp_stock_category 实测不可用 (无 stock_code 列), 勿用。 diff --git a/ddl_pms_v1.sql b/ddl_pms_v1.sql index 8434bcb..c09a110 100644 --- a/ddl_pms_v1.sql +++ b/ddl_pms_v1.sql @@ -4,6 +4,8 @@ -- 字符集: utf8mb4; 代码格式: Tushare 点式 (600000.SH); 时区: Asia/Shanghai -- 对应设计: POSITION_MGMT_DESIGN.md V0.4 §11 (1~10 号表) -- QMT_WS_PROTOCOL.md V1.0 (11~13 号表, 2026-07-28 追加) +-- QMT_WS_PROTOCOL.md §5.5 (14 号表, 2026-07-29 追加) +-- UPSTREAM_PLAN_API.md §9 (15 号表, 2026-07-31 追加) -- 注: pms_order_request (原表轮询通道) 已作废 —— 双方改定 WebSocket 直连, 见文件尾部 -- 三张 ws 通道表。该表本就未在此文件建立, 无需处理。 -- ===================================================================== @@ -302,3 +304,32 @@ CREATE TABLE IF NOT EXISTS pms_cash_flow ( UNIQUE KEY uk_kind_trade (kind, trade_no), KEY idx_ymd_kind (ymd, kind) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='现金流水 (费用不入成本, 只入现金账)'; + +-- 15. 上游计划名册快照 (2026-07-31 追加; UPSTREAM_PLAN_API.md §9) +-- 为什么需要它: 上游应答里的 `changes` 字段恒为 NULL (§4 Q4 已问未答), 而「谁新进榜、 +-- 谁掉榜、谁降了档」是上游观点变化最直接的信号。与其等对方补, 不如自己留底 —— 每次拉 +-- 计划落一份名册, 差异由 app/core/plan_diff.py 纯逻辑算。 +-- 顺带补上 §7.3 那个洞: `/plan` 是实时算的, 同一个 date 可以对应多个版本 (07-31 实测 +-- 相隔一小时的两次请求, 打分池 972/81 → 423/502、榜首换人), 而应答里没有 generated_at, +-- **下游此前没有任何办法判断"手上这份是哪一版"**。落了快照, 同一 plan_date 底下有几行 +-- 就是上游那天重算过几版, 复盘不用再靠猜。 +-- digest **不含 score**: /plan 实时算, score 末位天天抖; 算进指纹会让本表白涨而信息量 +-- 为零。含 rank 就够 —— 分数抖到改变次序才算真的变了, 没改次序对候选池毫无影响。 +-- **digest 上没有唯一键, 这是故意的。** 去重只跟"最新那行"比 (见 repo 层 +-- insert_plan_snapshot 的注释): 唯一键会让上游 A→B→A 改回去时第三次的 A 落不下去, +-- 库里最新仍是 B, 于是页面把 B→A 这次真实回退显示成 A→B —— 方向整个反了。 +CREATE TABLE IF NOT EXISTS pms_plan_snapshot ( + id BIGINT PRIMARY KEY AUTO_INCREMENT, + plan_date VARCHAR(16) NOT NULL COMMENT '上游 date (数据日, 通常是上一交易日)', + digest VARCHAR(40) NOT NULL COMMENT '名册指纹 sha1 (不含 score); 与最新一行同值即跳过', + fetched_ymd INT NOT NULL COMMENT 'PMS 拉到它的那天 YYYYMMDD', + fetched_at DATETIME NOT NULL, + n_main INT NOT NULL DEFAULT 0, + n_observe INT NOT NULL DEFAULT 0, + capped_main TINYINT NOT NULL DEFAULT 0 COMMENT '1 = 主榜吃满了请求的 top (尾部不可信)', + capped_obs TINYINT NOT NULL DEFAULT 0 COMMENT '1 = 观察档吃满了 obs_top; 与主榜各判各的', + roster_json MEDIUMTEXT NOT NULL COMMENT '名册 [{c,n,r,s,t,h,b}], 代码序', + meta_json TEXT NULL COMMENT 'counts / returned / requested / theme_cap 等', + KEY idx_date_id (plan_date, id), + KEY idx_fetched (fetched_at) +) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='上游计划名册快照 (榜单变化比对的底本)'; diff --git a/scripts/check_db.py b/scripts/check_db.py index acfc85a..1ce7a5f 100644 --- a/scripts/check_db.py +++ b/scripts/check_db.py @@ -6,7 +6,7 @@ 检查项: 1. 三个库连通性 (153 代理 / 因子库 / 指数库) - 2. pms_* 十三张表是否存在与当前行数 + 2. pms_* 十五张表是否存在与当前行数 3. 下游三表可读性 + 持仓数量列探测结果 (回填 QMT_INTERFACE_REQUIREMENTS A1/D1 用) 4. 行情 Redis (db13) 连通性与样本键 5. 运行参数表当前生效值 @@ -22,7 +22,11 @@ PMS_TABLES = ["pms_command", "pms_plan", "pms_position", "pms_lot", "pms_instruc "pms_proposal", "pms_action_ledger", "pms_daily_report", "pms_industry_map", "pms_runtime_param", # ws 直连通道三表 (协议 QMT_WS_PROTOCOL.md V1.0) - "pms_qmt_order", "pms_qmt_inbox", "pms_ws_state"] + "pms_qmt_order", "pms_qmt_inbox", "pms_ws_state", + # 14/15 号表。**加表时记得往这里加一行** —— 漏了的话 init_db 建了表、 + # check_db 却一声不吭, 表结构自检就有个静默的盲区。 + # (pms_cash_flow 就漏过一轮, 2026-07-31 补上) + "pms_cash_flow", "pms_plan_snapshot"] DOWNSTREAM = ["trading_position", "trading_order", "trading_buy_plan"] FAILED, WARNED = [], [] diff --git a/scripts/probe_plan_api.py b/scripts/probe_plan_api.py index e9576af..565d44f 100644 --- a/scripts/probe_plan_api.py +++ b/scripts/probe_plan_api.py @@ -6,11 +6,14 @@ docker compose run --rm --no-deps pms-web python scripts/probe_plan_api.py --date 2026-07-29 docker compose run --rm --no-deps pms-web python scripts/probe_plan_api.py --top 30 --with-price -干什么: 拿真实应答验证四件事 —— ① 接口通不通、② 字段口径与单测 fixture 是否一致、 -③ 按当前参数筛出来的候选池长什么样、④ (--with-price) 这些票行情里到底有没有价。 +干什么: 拿真实应答验证五件事 —— ① 接口通不通、② 字段口径与单测 fixture 是否一致、 +③ 按当前参数筛出来的候选池长什么样、④ (--with-price) 这些票行情里到底有没有价、 +⑤ 跟库里最新那份榜比, 谁新进谁掉榜、持仓票有没有被摘。 第 ④ 项最容易翻车: 计划不带价格, 价格取不到的票会在候选池里被静默剔除。 只读: 不会 upsert pms_industry_map (那由 /api/ops/plan-refresh 或盘前调度做), 也不下单。 +第 ⑤ 项默认也是只读 (拿手上这份跟库里最新快照比, 不写库); **唯一的写操作**是显式加 +`--snapshot` 时落一行 pms_plan_snapshot。 """ from __future__ import annotations @@ -166,6 +169,8 @@ def main(): ap.add_argument("--try-limit", action="store_true", help="探取全量的参数名 (上游默认只回主榜 20 条, counts 却是 961)") ap.add_argument("--limit-value", default="1000", help="--try-limit 用的取值 (默认 1000)") + ap.add_argument("--snapshot", action="store_true", + help="把这份名册落进 pms_plan_snapshot (**本脚本唯一的写操作**, 默认不写)") ap.add_argument("--json", action="store_true", help="原样打印解析后的结构") args = ap.parse_args() @@ -271,8 +276,53 @@ def main(): sorted(st.items(), key=lambda y: -y[1]))) print(" " + ", ".join(x["ts_code"] for x in sel["items"])) + # --- [7] 榜单变化 (PMS 自算; 上游 changes 恒为 null, 见 UPSTREAM_PLAN_API §9) + # 默认**只读**: 拿刚取到的这份跟库里最新快照比, 一个字都不写 —— 本脚本的契约 + # 是"只读探活", 破了它比少个功能糟得多。要落库请显式加 --snapshot。 + print("[7] 榜单变化" + (" (落库)" if args.snapshot else " (只读, 不写库)")) + try: + # **必须先比后写。** 反过来的话 preview 会拿这份计划跟"刚写进去的它自己"比, + # 于是永远报"完全相同、无变化" —— 加个 --snapshot 就把要看的东西看没了。 + ch = pf.preview_changes(plan) + if args.snapshot: + snap = pf.snapshot(plan, keep=ps.get_int("PMS_PLAN_SNAPSHOT_KEEP", 200)) + print(f" 落快照: " + f"{'新版本已入库' if snap.get('stored') else '与库里最新那份相同, 未新增'} " + f"(digest={snap.get('digest', '')[:12]} {snap.get('rows')} 只)") + d = ch.get("diff") + if not ch.get("ok"): + print(f" {ch.get('hint') or '算不了'}") + else: + if d.get("prev_broken"): + print(" ! 上一版名册**读不出来** (快照损坏) —— 这次比不了, " + "下面的空白不代表没有变化") + elif ch.get("same_as_stored"): + print(" 与库里最新那份**完全相同** (指纹一致) —— 上游没重算") + kind = {"first": "库里还没有快照, 无可比对的上一版 (跑一次 --snapshot 或页面强刷即有)", + "cross_day": "跨交易日的正常升降档", + "same_date_revision": "**同一计划日的新版本 —— 上游重算过** (§7.3 已知行为)"} + print(f" {kind.get(d['kind'], d['kind'])}") + print(f" {d['note']}") + h = d.get("held") or {} + for x in (h.get("exited") or []): + print(f" ! 持仓票掉榜: {x['ts_code']} {x['name']} " + f"(上版 rank {x['rank']} / {x['tier']})") + for x in (h.get("tier_down") or []): + print(f" ! 持仓票降档: {x['ts_code']} {x['name']} " + f"{x['tier_from']} → {x['tier_to']}") + for x in (h.get("coverage_lost") or []): + print(f" ! 持仓票丢了券商覆盖 (主榜→观察档): {x['ts_code']} {x['name']}") + if not h.get("n_watch") and d["kind"] != "first" and not d.get("prev_broken"): + print(" 持仓票没被摘、没降档、也没丢券商覆盖") + log = pf.snapshot_log(limit=30) + if log.get("revised"): + print(" 同一计划日出现多版: " + + ", ".join(f"{k}×{v}" for k, v in log["revised"].items())) + except Exception as e: + print(f" ! 算榜单变化失败: {type(e).__name__}: {e}") + if args.json: - print("[7] 解析结构") + print("[8] 解析结构") print(json.dumps({k: v for k, v in plan.items() if k not in ("main", "observe")}, ensure_ascii=False, indent=2)) return 0 diff --git a/scripts/reset_ledger.py b/scripts/reset_ledger.py index 3fdd465..4e476c2 100644 --- a/scripts/reset_ledger.py +++ b/scripts/reset_ledger.py @@ -57,9 +57,11 @@ COMMAND_TABLES = ["pms_plan", "pms_command", "pms_proposal"] # ws 通道两表 —— 默认**不清**, 见模块注释第 1 条; --purge-channel 才清, 且必须成对 CHANNEL_TABLES = ["pms_qmt_order", "pms_qmt_inbox"] # 绝不碰: 下游三表 / 行业映射 / 参数表 (含页面调过的全部业务参数与信号当日去重集合) +# pms_plan_snapshot 也在此列: 它记的是**上游给过我们什么**, 与本方账本无关。清账是"我方 +# 重来一次", 不该把上游的历史一起抹掉 —— 那反而让"这份榜是哪天哪一版"更查不清。 NEVER_TOUCH = ["trading_order", "trading_position", "trading_buy_plan", "strategy_daily_results", "gp_stock_category", - "pms_industry_map", "pms_runtime_param"] + "pms_industry_map", "pms_runtime_param", "pms_plan_snapshot"] def _count(fetch_one, t): diff --git a/scripts/run_tests.py b/scripts/run_tests.py index 9be510f..a75cfd8 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -12,8 +12,9 @@ test_batch5_units.py 决策系统信号流解析与消化口径 (8 例) test_batch6_units.py ws 通道: 测试向量/签名/公钥/水位/DDL/逐笔入账 (65 例) test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 (32 例) + test_batch8_units.py 榜单变化: 名册指纹/三种语义/尾部闸/落库往返 (60 例) test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (58 例) - 共 244 例 + 共 304 例 任一子集失败即整体失败 (退出码 1)。 """ import os @@ -24,7 +25,7 @@ HERE = os.path.dirname(os.path.abspath(__file__)) ROOT = os.path.dirname(HERE) SUITES = ["test_core_units.py", "test_batch2_units.py", "test_batch3_units.py", "test_batch4_units.py", "test_batch5_units.py", "test_batch6_units.py", - "test_batch7_units.py", "test_wiring.py"] + "test_batch7_units.py", "test_batch8_units.py", "test_wiring.py"] def main(): diff --git a/scripts/test_batch6_units.py b/scripts/test_batch6_units.py index 4462cbf..95f6197 100644 --- a/scripts/test_batch6_units.py +++ b/scripts/test_batch6_units.py @@ -557,7 +557,7 @@ def run(): "DEFAULT 'NONE' COMMENT 'NONE/REQUESTED/SENT'"): eq(find_adjacent_literals(good), [], f"误报: {good[:40]}") - @case("DDL 文件本身体检通过 (14 张表 + 1 条初始行)") + @case("DDL 文件本身体检通过 (15 张表 + 1 条初始行)") def _(): sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from init_db import DDL_FILE, find_adjacent_literals, parse_statements @@ -567,7 +567,9 @@ def run(): eq(broken, [], "有认不出的残句, 说明 DDL 切分被破坏") for _, tbl, s in stmts: eq(find_adjacent_literals(s), [], f"{tbl} 有相邻字面量") - eq(len([1 for k, _, _ in stmts if k == "table"]), 14) + # 加表就要来这里 +1 —— 这一行是 DDL 与代码之间唯一的哨兵, 它不动就说明新表没进 + # ddl_pms_v1.sql (init_db 只认这个文件, 建不出来的表在实机上才会报"表不存在") + eq(len([1 for k, _, _ in stmts if k == "table"]), 15) eq(len([1 for k, _, _ in stmts if k == "seed"]), 1) @case("通道三表的 SQL 全部单表合规") diff --git a/scripts/test_batch8_units.py b/scripts/test_batch8_units.py new file mode 100644 index 0000000..ecff51d --- /dev/null +++ b/scripts/test_batch8_units.py @@ -0,0 +1,763 @@ +# -*- coding: utf-8 -*- +""" +第八批模块单测 (零外部依赖, 不连库不触网) +========================================== +运行: 在 tradingSystem 仓库根目录执行 python scripts/test_batch8_units.py +覆盖: 上游榜单的版本比对 (app/core/plan_diff.py) —— 名册与指纹、三种比对语义、 + 进出榜的尾部闸、档位升降与券商覆盖翻转、名次跳变、持仓票视图。 + +这一批的重点不是"能不能算出差异", 而是**不报假变化**: 榜单是按 top 截断过的, 尾部的进出 +是截断噪音而不是上游观点变化。假警报天天响, 提示就等于没有 —— 与 `_capped` 治的是同一个病。 +""" +import os +import sys +import traceback + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +from app.core import plan_diff as pd # noqa: E402 + +RESULTS = [] + + +def case(name): + def deco(fn): + RESULTS.append((name, fn)) + return fn + return deco + + +# ================================================================ fixture +def _r(code, rank, score, tier="强传导", theme="储能", bucket="main", name=None): + """名册单行 (plan_diff 内部形态)。""" + return {"c": code, "n": name or code[:6], "r": rank, "s": score, + "t": tier, "h": theme, "b": bucket} + + +def roster(*rows) -> dict: + return {r["c"]: r for r in rows} + + +def plan(main=(), observe=(), date="2026-07-30", capped_main=False, capped_obs=False) -> dict: + """plan_feed.parse_plan() 结果的最小形态。 + + `returned` / `truncated` 也要带上 —— 快照落库时靠它们记条数与截断标志, 少了的话 + "上一版坏没坏"这类判据在单测里恒为假, 用例看着过了实际什么都没验。 + """ + def row(x, bucket): + return {"ts_code": x["c"], "name": x["n"], "rank": x["r"], "score": x["s"], + "tier": x["t"], "theme": x["h"], "bucket": bucket} + return {"date": date, + "main": [row(x, "main") for x in main], + "observe": [row(x, "observe") for x in observe], + "returned": {"main": len(main), "observe": len(observe)}, + "truncated": {"main": bool(capped_main), "observe": bool(capped_obs)}} + + +# ================================================================ [A] 名册与指纹 +@case("[A1] roster_of 从计划结构取出名册") +def t_a1(): + ro = pd.roster_of(plan(main=[_r("600418.SH", 1, 242.2)], + observe=[_r("600877.SH", 1, 101.1, tier=None, bucket="observe")])) + assert set(ro) == {"600418.SH", "600877.SH"} + assert ro["600418.SH"]["b"] == "main" and ro["600877.SH"]["b"] == "observe" + assert ro["600418.SH"]["t"] == "强传导" and ro["600877.SH"]["t"] is None + + +@case("[A2] 同一只票两档都出现时以主榜为准") +def t_a2(): + p = plan(main=[_r("600418.SH", 3, 240.0)], + observe=[_r("600418.SH", 1, 99.0, tier=None, bucket="observe")]) + ro = pd.roster_of(p) + assert len(ro) == 1 + assert ro["600418.SH"]["b"] == "main", "主榜在前, 观察档那条不该覆盖它" + + +@case("[A3] 指纹稳定: 同一名册两次算结果一致") +def t_a3(): + a = roster(_r("600000.SH", 1, 10.0), _r("600418.SH", 2, 9.0)) + b = roster(_r("600418.SH", 2, 9.0), _r("600000.SH", 1, 10.0)) # 插入次序不同 + assert pd.digest_of(a) == pd.digest_of(b), "指纹不能依赖 dict 的插入次序" + + +@case("[A4] 指纹不含 score —— 分数抖动但次序没变, 不算新版本") +def t_a4(): + a = roster(_r("600000.SH", 1, 242.2801)) + b = roster(_r("600000.SH", 1, 242.2799)) + assert pd.digest_of(a) == pd.digest_of(b), ( + "/plan 是实时算的, score 末位天天抖; 算进指纹会让快照表白涨而信息量为零") + + +@case("[A5] 指纹含 rank —— 次序变了就是新版本") +def t_a5(): + assert pd.digest_of(roster(_r("600000.SH", 1, 10.0))) != \ + pd.digest_of(roster(_r("600000.SH", 2, 10.0))) + + +@case("[A6] 指纹含 tier / theme / bucket") +def t_a6(): + base = roster(_r("600000.SH", 1, 10.0)) + assert pd.digest_of(base) != pd.digest_of(roster(_r("600000.SH", 1, 10.0, tier="弱传导"))) + assert pd.digest_of(base) != pd.digest_of(roster(_r("600000.SH", 1, 10.0, theme="整机制造"))) + assert pd.digest_of(base) != pd.digest_of(roster(_r("600000.SH", 1, 10.0, bucket="observe"))) + + +@case("[A7] 名册 ↔ 落库行 往返不丢字段") +def t_a7(): + ro = roster(_r("600418.SH", 7, 233.3, tier="弱传导", theme="整车")) + back = pd.roster_from_rows(pd.roster_rows(ro)) + assert back == ro + assert pd.digest_of(back) == pd.digest_of(ro) + + +@case("[A8] 落库行里的坏数据跳过而不是整份废掉") +def t_a8(): + back = pd.roster_from_rows([{"c": "600000.SH", "r": 1}, "坏行", None, {}, {"r": 5}]) + assert list(back) == ["600000.SH"] + + +@case("[A9] 名册行按代码序输出 —— 落库文本稳定") +def t_a9(): + rows = pd.roster_rows(roster(_r("600418.SH", 2, 9.0), _r("600000.SH", 1, 10.0))) + assert [r["c"] for r in rows] == ["600000.SH", "600418.SH"] + + +# ================================================================ [B] 档位序 +@case("[B1] 档位强弱: 强传导 > 弱传导 > 无传导") +def t_b1(): + assert pd.tier_rank("强传导") > pd.tier_rank("弱传导") > pd.tier_rank("无传导") + + +@case("[B2] 未知档位返回 None (判不了就不判)") +def t_b2(): + assert pd.tier_rank("超强传导") is None and pd.tier_rank(None) is None + + +@case("[B3] 未知档位不参与升降判定 —— 宁可不报, 不能报反") +def t_b3(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0, tier="超强传导")), + roster(_r("600000.SH", 1, 10.0, tier="强传导")), + prev_date="2026-07-29", curr_date="2026-07-30") + assert d["counts"]["tier_up"] == 0 and d["counts"]["tier_down"] == 0 + + +# ================================================================ [C] 首次 +@case("[C1] 没有上一版时一条变化都不报") +def t_c1(): + d = pd.diff(None, roster(*[_r(f"60000{i}.SH", i, 100 - i) for i in range(1, 9)]), + curr_date="2026-07-30") + assert d["kind"] == pd.KIND_FIRST + assert sum(d["counts"].values()) == 0, "首次落快照若报'新进榜', 会是整整一榜的假变化" + assert d["roster_size"]["curr"] == 8 + + +@case("[C2] 首次的 note 说明从下一次开始才有比对") +def t_c2(): + d = pd.diff({}, roster(_r("600000.SH", 1, 10.0))) + assert "首次" in d["note"] and d["prev_broken"] is False + + +@case("[C3] 上一版读不出来要说成故障, 不能长得跟'首次'一样") +def t_c3(): + # 快照的 roster_json 坏了 → 解析成空。若与"首次"同形, 一条真实变化都不报却显示正常。 + d = pd.diff({}, roster(_r("600000.SH", 1, 10.0)), + prev_date="2026-07-29", curr_date="2026-07-30", prev_broken=True) + assert d["prev_broken"] is True + assert "读不出来" in d["note"] and "比不了" in d["note"] + assert "首次" not in d["note"] + + +@case("[C4] 首次也要回显调用方给的参数, 不能写死默认值") +def t_c4(): + d = pd.diff(None, roster(_r("600000.SH", 1, 10.0)), rank_jump=80, tail_guard=0.2) + assert d["params"] == {"rank_jump": 80, "tail_guard": 0.2}, "写死会让页面显示错的口径" + + +# ================================================================ [D] 三种比对语义 +@case("[D1] 计划日不同 = 跨日的正常升降档") +def t_d1(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0)), roster(_r("600000.SH", 1, 10.0)), + prev_date="2026-07-29", curr_date="2026-07-30") + assert d["kind"] == pd.KIND_CROSS_DAY + + +@case("[D2] 同一计划日的不同版本要单独成一类") +def t_d2(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0)), roster(_r("600000.SH", 2, 10.0)), + prev_date="2026-07-30", curr_date="2026-07-30") + assert d["kind"] == pd.KIND_REVISION + + +@case("[D3] 同日新版本的结论要点明'候选池与此前不是同一份'") +def t_d3(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0)), roster(_r("600418.SH", 1, 11.0)), + prev_date="2026-07-30", curr_date="2026-07-30") + assert "重算" in d["note"] and "候选池" in d["note"] + + +# ================================================================ [E] 进出榜与尾部闸 +@case("[E1] 新进榜与掉榜各归各位") +def t_e1(): + prev = roster(_r("600000.SH", 1, 10.0), _r("600418.SH", 2, 9.0)) + curr = roster(_r("600000.SH", 1, 10.0), _r("600519.SH", 2, 9.5)) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30") + assert [x["ts_code"] for x in d["entered"]] == ["600519.SH"] + assert [x["ts_code"] for x in d["exited"]] == ["600418.SH"] + + +@case("[E2] 掉榜那条带的是**上一版**的排名与档位") +def t_e2(): + prev = roster(_r("600418.SH", 7, 233.3, tier="弱传导")) + d = pd.diff(prev, roster(_r("600000.SH", 1, 10.0)), + prev_date="2026-07-29", curr_date="2026-07-30") + x = d["exited"][0] + assert x["rank"] == 7 and x["tier"] == "弱传导", "掉榜的票在这一版没有数据, 只能示上一版" + + +@case("[E3] 两版都没吃满 top 时不设尾部闸 —— 消失就是真消失") +def t_e3(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 10)]) # 第 10 名没了 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + prev_capped=False, curr_capped=False) + assert d["counts"]["exited"] == 1 and d["tail_churn"]["exited"] == 0 + # tail_guarded 报的是各档闸位在第几名, None = 该档没设闸 (让"为什么这条没报"可查) + assert d["tail_guarded"] == {"entered": {"main": None, "observe": None}, + "exited": {"main": None, "observe": None}} + + +@case("[E4] 这一版吃满 top 时, 榜尾掉榜归 tail_churn 不列名") +def t_e4(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 10)]) # 原第 10 名没了 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + curr_capped=True) # 榜长 9, 闸位 4.5, 上一版 rank 10 在闸外 + assert d["counts"]["exited"] == 0 + assert d["tail_churn"]["exited"] == 1, "它可能只是掉到我们要的条数之外, 不是上游摘了它" + + +@case("[E5] 吃满 top 也照报榜首掉榜") +def t_e5(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(2, 11)]) # 榜首没了 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", curr_capped=True) + assert [x["ts_code"] for x in d["exited"]] == ["600001.SH"] + assert d["tail_churn"]["exited"] == 0 + + +@case("[E6] 新进榜的尾部闸看的是**上一版**吃没吃满") +def t_e6(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = dict(prev) + curr["600099.SH"] = _r("600099.SH", 9, 91.5) # 新面孔排在榜尾 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", prev_capped=True) + assert d["counts"]["entered"] == 0 + assert d["tail_churn"]["entered"] == 1, "上一版吃满了, 这张脸当时可能就在榜上只是没露面" + + +@case("[E7] 上一版吃满时, 榜首的新面孔照报") +def t_e7(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = dict(prev) + curr["600099.SH"] = _r("600099.SH", 1, 200.0) # 新面孔直接冲到第一 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", prev_capped=True) + assert [x["ts_code"] for x in d["entered"]] == ["600099.SH"] + + +@case("[E8] 两道闸互不干扰: 这一版截断不该影响新进榜的判定") +def t_e8(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = dict(prev) + curr["600099.SH"] = _r("600099.SH", 9, 91.5) + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + prev_capped=False, curr_capped=True) + assert d["counts"]["entered"] == 1, "上一版没吃满, 那这张新面孔就是真的新" + assert d["tail_churn"]["entered"] == 0 + + +@case("[E9] 整档消失不被尾部闸吞掉") +def t_e9(): + prev = roster(_r("600877.SH", 1, 101.0, tier=None, bucket="observe"), + _r("600000.SH", 1, 10.0)) + curr = roster(_r("600000.SH", 1, 10.0)) # 观察档整档没了 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", curr_capped=True) + assert [x["ts_code"] for x in d["exited"]] == ["600877.SH"], ( + "这一版观察档一条都没有, 谈不上被 top 截断 —— 整档消失本身就是要报的大事") + + +@case("[E10] 尾部闸按档各算各的 (闸位取自各自那一档的长度)") +def t_e10(): + prev = roster(*([_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)] + + [_r("600877.SH", 1, 101.0, tier=None, bucket="observe")])) + curr = roster(*([_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 10)] + + [_r(f"6009{i:02d}.SH", i, 90 - i, tier=None, bucket="observe") + for i in range(1, 11)])) + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + curr_capped={"main": True, "observe": True}) + # 主榜榜尾那只走 churn; 观察档那只 rank=1、这一版观察档有 10 条(闸位 5), 在闸内 → 照报 + assert [x["ts_code"] for x in d["exited"]] == ["600877.SH"] + assert d["tail_churn"]["exited"] == 1 + + +@case("[E11] 主榜吃满不代表观察档也吃满 —— 别拿主榜的截断解释观察档的摘牌") +def t_e11(): + # 上游两个独立参数: top=10 吃满了, obs_top=100 只回了 4 条 (没吃满)。 + # 观察档那只真被摘了 —— 若两档共用主榜那一个标志, 它会被记成"截断噪音"永远看不见。 + prev = roster(*([_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)] + + [_r(f"6009{i:02d}.SH", i, 90 - i, tier=None, bucket="observe") + for i in range(1, 5)])) + curr = roster(*([_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)] + + [_r(f"6009{i:02d}.SH", i, 90 - i, tier=None, bucket="observe") + for i in range(1, 4)])) # 观察档第 4 名被摘 + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + curr_capped={"main": True, "observe": False}) + assert [x["ts_code"] for x in d["exited"]] == ["600904.SH"], ( + "观察档没吃满就不该设闸 —— 拿主榜的截断解释观察档, 会吞掉一次真实摘牌") + assert d["tail_churn"]["exited"] == 0 + # 反过来: 两档都按吃满算, 这条就会被吞掉 (证明上面那条断言不是白写的) + d2 = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", curr_capped=True) + assert d2["counts"]["exited"] == 0 and d2["tail_churn"]["exited"] == 1 + + +@case("[E12] 持仓票不受榜尾闸约束 —— 漏报一条就是一个仓位") +def t_e12(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 10)]) + # 非持仓: 上一版 rank 10 在闸外 → 吞掉 + d0 = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", curr_capped=True) + assert d0["counts"]["exited"] == 0 + # 同一只票是持仓: 照报, 且记一笔豁免 + d1 = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + curr_capped=True, held=["600010.SH"]) + assert [x["ts_code"] for x in d1["exited"]] == ["600010.SH"] + assert d1["tail_churn"]["exited"] == 0 and d1["tail_churn"]["held_exempt"] == 1 + assert d1["held"]["n_watch"] == 1 + + +@case("[E13] rank 判不出来时照报, 不当噪音吞掉") +def t_e13(): + prev = roster(_r("600000.SH", None, 10.0), _r("600001.SH", 1, 99.0)) + curr = roster(_r("600001.SH", 1, 99.0)) + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", curr_capped=True) + assert [x["ts_code"] for x in d["exited"]] == ["600000.SH"], ( + "这个函数的'是'代表吞掉一条变化, 所以拿不准时必须放行") + + +# ================================================================ [F] 档位与券商覆盖 +@case("[F1] 升档与降档") +def t_f1(): + prev = roster(_r("600000.SH", 1, 10.0, tier="弱传导"), _r("600418.SH", 2, 9.0, tier="强传导")) + curr = roster(_r("600000.SH", 1, 10.0, tier="强传导"), _r("600418.SH", 2, 9.0, tier="弱传导")) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30") + assert [x["ts_code"] for x in d["tier_up"]] == ["600000.SH"] + assert d["tier_up"][0]["tier_from"] == "弱传导" and d["tier_up"][0]["tier_to"] == "强传导" + assert [x["ts_code"] for x in d["tier_down"]] == ["600418.SH"] + + +@case("[F2] 主榜 ↔ 观察档 互换 = 券商覆盖翻转, 单独一段") +def t_f2(): + prev = roster(_r("600000.SH", 1, 10.0)) + curr = roster(_r("600000.SH", 1, 10.0, tier=None, bucket="observe")) + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30") + assert len(d["bucket_moved"]) == 1 + m = d["bucket_moved"][0] + assert m["moved_from"] == "main" and m["moved_to"] == "observe" + assert d["counts"]["exited"] == 0, "换档不是掉榜 —— 票还在, 只是没券商覆盖了" + + +@case("[F3] 换了档就不算名次跳变 —— 两档的 rank 不可比") +def t_f3(): + prev = roster(_r("600000.SH", 300, 10.0)) + curr = roster(_r("600000.SH", 1, 10.0, tier=None, bucket="observe")) + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", rank_jump=10) + assert d["counts"]["rank_jump"] == 0 + assert d["counts"]["bucket_moved"] == 1 + + +# ================================================================ [G] 名次跳变 +@case("[G1] 跳变到阈值才报") +def t_g1(): + prev = roster(_r("600000.SH", 100, 10.0), _r("600418.SH", 20, 20.0)) + curr = roster(_r("600000.SH", 20, 30.0), _r("600418.SH", 15, 22.0)) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30", rank_jump=50) + assert [x["ts_code"] for x in d["rank_jump"]] == ["600000.SH"] + assert d["rank_jump"][0]["rank_from"] == 100 and d["rank_jump"][0]["rank_delta"] == -80 + + +@case("[G2] 按跳变幅度降序") +def t_g2(): + prev = roster(_r("600000.SH", 200, 1.0), _r("600418.SH", 60, 2.0)) + curr = roster(_r("600000.SH", 100, 5.0), _r("600418.SH", 1, 9.0)) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30", rank_jump=50) + assert [x["ts_code"] for x in d["rank_jump"]] == ["600000.SH", "600418.SH"] + + +@case("[G3] rank 缺失时不算跳变") +def t_g3(): + prev = roster(_r("600000.SH", None, 1.0)) + curr = roster(_r("600000.SH", 1, 9.0)) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30", rank_jump=1) + assert d["counts"]["rank_jump"] == 0 + + +# ================================================================ [H] 持仓票视图 +@case("[H1] 持仓票掉榜进 n_watch") +def t_h1(): + prev = roster(_r("600000.SH", 1, 10.0), _r("600418.SH", 2, 9.0)) + curr = roster(_r("600418.SH", 1, 9.0)) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30", + held=["600000.SH"]) + assert [x["ts_code"] for x in d["held"]["exited"]] == ["600000.SH"] + assert d["held"]["n_watch"] == 1 + assert d["exited"][0]["held"] is True + + +@case("[H2] 持仓票降档也进 n_watch") +def t_h2(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0, tier="强传导")), + roster(_r("600000.SH", 1, 10.0, tier="弱传导")), + prev_date="2026-07-29", curr_date="2026-07-30", held=["600000.SH"]) + assert d["held"]["n_watch"] == 1 and len(d["held"]["tier_down"]) == 1 + + +@case("[H2b] 持仓票掉进观察档 = 丢了券商覆盖, 也要进 n_watch") +def t_h2b(): + # 观察档的行没有 tier, 所以这既不算 exited 也不算 tier_down —— 不显式算进来的话, + # 一只持仓票丢了估值锚, 页面红条不亮、探活脚本还会打一句"没被摘也没降档"。 + d = pd.diff(roster(_r("600000.SH", 3, 240.0, tier="强传导")), + roster(_r("600000.SH", 5, 99.0, tier=None, bucket="observe")), + prev_date="2026-07-29", curr_date="2026-07-30", held=["600000.SH"]) + assert d["counts"]["exited"] == 0 and d["counts"]["tier_down"] == 0 + assert [x["ts_code"] for x in d["held"]["coverage_lost"]] == ["600000.SH"] + assert d["held"]["n_watch"] == 1 + assert "券商覆盖" in d["note"] + + +@case("[H2c] 掉进主榜(拿到券商覆盖)不算要盯的事") +def t_h2c(): + d = pd.diff(roster(_r("600000.SH", 5, 99.0, tier=None, bucket="observe")), + roster(_r("600000.SH", 3, 240.0, tier="强传导")), + prev_date="2026-07-29", curr_date="2026-07-30", held=["600000.SH"]) + assert d["counts"]["bucket_moved"] == 1 + assert d["held"]["coverage_lost"] == [] and d["held"]["n_watch"] == 0 + + +@case("[H3] 持仓票升档只进 n_total, 不进 n_watch") +def t_h3(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0, tier="弱传导")), + roster(_r("600000.SH", 1, 10.0, tier="强传导")), + prev_date="2026-07-29", curr_date="2026-07-30", held=["600000.SH"]) + assert d["held"]["n_watch"] == 0 and d["held"]["n_total"] == 1 + + +@case("[H4] 非持仓票不进持仓视图") +def t_h4(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0)), roster(_r("600418.SH", 1, 10.0)), + prev_date="2026-07-29", curr_date="2026-07-30", held=["600519.SH"]) + assert d["held"]["n_total"] == 0 and d["counts"]["exited"] == 1 + + +@case("[H5] 持仓票掉榜要在结论句里显眼") +def t_h5(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0)), roster(_r("600418.SH", 1, 10.0)), + prev_date="2026-07-29", curr_date="2026-07-30", held=["600000.SH"]) + assert "持仓票" in d["note"] and "掉榜" in d["note"] + + +# ================================================================ [I] 健壮性 +@case("[I1] 空对空不炸") +def t_i1(): + d = pd.diff({}, {}, prev_date="2026-07-30", curr_date="2026-07-30") + assert d["kind"] == pd.KIND_FIRST and sum(d["counts"].values()) == 0 + + +@case("[I2] 名册全没了也报得出来") +def t_i2(): + d = pd.diff(roster(_r("600000.SH", 1, 10.0)), {}, + prev_date="2026-07-29", curr_date="2026-07-30") + assert d["counts"]["exited"] == 1 and d["roster_size"]["curr"] == 0 + + +@case("[I3] tail_guard 越界值被夹回 [0,1]") +def t_i3(): + prev = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11)]) + curr = roster(*[_r(f"6000{i:02d}.SH", i, 100 - i) for i in range(1, 11) if i != 5]) + for given, clamped, reported in ((5.0, 1.0, 1), (-1.0, 0.0, 0)): + d = pd.diff(prev, curr, prev_date="2026-07-30", curr_date="2026-07-30", + curr_capped=True, tail_guard=given) + assert d["params"]["tail_guard"] == clamped + # 闸开到底(1.0): 榜内的第 5 名掉了要报; 闸关到底(0.0): 一律当截断噪音 + assert d["counts"]["exited"] == reported + + +@case("[I4] counts 与各段长度始终一致") +def t_i4(): + prev = roster(_r("600000.SH", 1, 10.0, tier="强传导"), _r("600418.SH", 200, 2.0)) + curr = roster(_r("600000.SH", 1, 10.0, tier="弱传导"), _r("600519.SH", 3, 8.0)) + d = pd.diff(prev, curr, prev_date="2026-07-29", curr_date="2026-07-30") + for k, n in d["counts"].items(): + assert len(d[k]) == n, f"{k} 的 counts 与实际条数对不上" + + +# ================================================================ [J] 服务层落库往返 +# 纯逻辑之外唯一容易藏 bug 的地方: 名册 → JSON 落库 → 读回来 → 比对 这条往返。 +# 单测不连库 —— 拿内存桩顶掉 pms_repo 与 param_store (plan_feed 是函数内 import, 装在 +# sys.modules 上就生效; 与 test_batch7 用假 requests 是同一套路)。 +def _with_fake_repo(params=None): + import json as _j + import types + store = {"rows": [], "seq": 0} + + repo = types.ModuleType("app.repo.pms_repo") + + def insert_plan_snapshot(*, plan_date, digest, roster, meta=None, n_main=0, n_observe=0, + capped_main=False, capped_obs=False, fetched_at=None): + # 与生产同一条判据: 只跟**最新那行**比, 不是"曾经出现过的任何一行"。 + # 早先桩里写的是后者 (模拟唯一键), 于是 A→B→A 这个 case 在桩上永远复现不出来 —— + # 桩比生产更严, 单测就成了给错实现背书。 + last = max(store["rows"], key=lambda x: x["id"]) if store["rows"] else None + if last and last["plan_date"] == plan_date and last["digest"] == digest: + return 0 + store["seq"] += 1 + store["rows"].append({ + "id": store["seq"], "plan_date": plan_date, "digest": digest, + "fetched_ymd": 20260731, "fetched_at": f"t{store['seq']}", + "n_main": n_main, "n_observe": n_observe, + "capped_main": bool(capped_main), "capped_obs": bool(capped_obs), + # 真表是 MEDIUMTEXT, 落进去的是 JSON 文本、读回来才转对象。桩必须照做 —— + # 否则"名册里塞了不可序列化的东西"这类事故单测里根本看不见。 + "roster_json": _j.dumps(roster, ensure_ascii=False), "meta": meta or {}}) + return 1 + + def latest_plan_snapshots(limit=2): + out = [] + for r in sorted(store["rows"], key=lambda x: -x["id"])[:limit]: + r2 = dict(r) + r2["roster"] = _j.loads(r2.pop("roster_json")) + out.append(r2) + return out + + def list_plan_snapshots(*, plan_date=None, limit=50): + rows = [dict(r) for r in sorted(store["rows"], key=lambda x: -x["id"])] + if plan_date: + rows = [r for r in rows if r["plan_date"] == plan_date] + for r in rows: + r.pop("roster_json", None) + return rows[:limit] + + repo.insert_plan_snapshot = insert_plan_snapshot + repo.latest_plan_snapshots = latest_plan_snapshots + repo.list_plan_snapshots = list_plan_snapshots + repo.prune_plan_snapshots = lambda keep=200: 0 + repo.list_positions = lambda only_open=False: [] + + ps = types.ModuleType("app.services.param_store") + vals = dict(params or {}) + ps.get = lambda k, d=None: vals.get(k, d) + ps.get_int = lambda k, d=0: int(vals.get(k, d)) + ps.get_float = lambda k, d=0.0: float(vals.get(k, d)) + ps.get_bool = lambda k, d=False: bool(vals.get(k, d)) + ps.get_list = lambda k, d=(): list(vals.get(k, d)) + + saved = {n: sys.modules.get(n) for n in ("app.repo.pms_repo", "app.services.param_store")} + sys.modules["app.repo.pms_repo"] = repo + sys.modules["app.services.param_store"] = ps + + def restore(): + for n, m in saved.items(): + if m is None: + sys.modules.pop(n, None) + else: + sys.modules[n] = m + return store, restore + + +@case("[J1] 落库往返: 名册经 JSON 存取后指纹不变") +def t_j1(): + import json as _j + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + p = plan(main=[_r("600418.SH", 1, 242.24, theme="整车"), + _r("600000.SH", 2, 241.55, tier="弱传导", theme=None)], + observe=[_r("600877.SH", 1, 101.1, tier=None, bucket="observe")], + date="2026-07-30") + out = pf.snapshot(p) + assert out["stored"] is True and out["rows"] == 3 + back = pd.roster_from_rows(_j.loads(store["rows"][0]["roster_json"])) + assert pd.digest_of(back) == out["digest"], "存取一轮指纹就变了, 比对全废" + assert back["600000.SH"]["h"] is None, "theme 为空的票也要原样存回来" + finally: + restore() + + +@case("[J2] 同一份榜重复落只留一行") +def t_j2(): + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + p = plan(main=[_r("600418.SH", 1, 242.24)], date="2026-07-30") + assert pf.snapshot(p)["stored"] is True + assert pf.snapshot(p)["stored"] is False, "内容没变还新增一行, 表会白涨" + assert len(store["rows"]) == 1 + finally: + restore() + + +@case("[J3] score 变了但次序没变, 不算新版本") +def t_j3(): + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + pf.snapshot(plan(main=[_r("600418.SH", 1, 242.2801)], date="2026-07-30")) + pf.snapshot(plan(main=[_r("600418.SH", 1, 242.2799)], date="2026-07-30")) + assert len(store["rows"]) == 1 + finally: + restore() + + +@case("[J4] 只读预览不写库, 且认得出'与库里那份相同'") +def t_j4(): + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + p = plan(main=[_r("600418.SH", 1, 242.24)], date="2026-07-30") + pf.snapshot(p) + n = len(store["rows"]) + out = pf.preview_changes(p) + assert out["ok"] and out["same_as_stored"] is True + assert len(store["rows"]) == n, "preview 必须只读 —— 探活脚本的只读契约靠它" + finally: + restore() + + +@case("[J5] 只读预览能比出手上这份与库里那份的差异") +def t_j5(): + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + pf.snapshot(plan(main=[_r("600418.SH", 1, 242.24)], date="2026-07-30")) + out = pf.preview_changes(plan(main=[_r("600519.SH", 1, 250.0)], date="2026-07-30")) + d = out["diff"] + assert d["kind"] == pd.KIND_REVISION + assert [x["ts_code"] for x in d["entered"]] == ["600519.SH"] + assert [x["ts_code"] for x in d["exited"]] == ["600418.SH"] + assert len(store["rows"]) == 1 + finally: + restore() + + +@case("[J6] 快照关掉时 changes 明说不可用, 而不是装作没变化") +def t_j6(): + from app.services import plan_feed as pf + _, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": False}) + try: + out = pf.changes() + assert out["ok"] is False and out["diff"] is None + assert "PMS_PLAN_SNAPSHOT" in out["hint"], "得说清是关掉了, 不是没有变化" + finally: + restore() + + +@case("[J7] 快照落库失败不阻断候选池, 但要把错带出来") +def t_j7(): + from app.services import plan_feed as pf + _, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + def _boom(**kw): + raise RuntimeError("DB 挂了") + sys.modules["app.repo.pms_repo"].insert_plan_snapshot = _boom + out = pf._snapshot_quiet(plan(main=[_r("600418.SH", 1, 1.0)], date="2026-07-30")) + assert out["stored"] is False and "DB 挂了" in out["error"] + finally: + restore() + + +@case("[J7b] 上游 A→B→A 改回去时, 比对方向不能反") +def t_j7b(): + # 曾经用 (plan_date, digest) 唯一键去重, 于是第三次的 A 落不下去、库里最新仍是 B, + # 页面把 B→A 这次真实回退显示成 A→B —— 方向整个反了, 还顺带诬告一句"落库失败"。 + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + a = plan(main=[_r("600418.SH", 1, 242.24)], date="2026-07-30") + b = plan(main=[_r("600519.SH", 1, 250.00)], date="2026-07-30") + pf.snapshot(a) + pf.snapshot(b) + assert pf.snapshot(a)["stored"] is True, "改回去也是一次真实变化, 必须落得下去" + assert len(store["rows"]) == 3 + out = pf.preview_changes(a) + assert out["same_as_stored"] is True + d = pf.changes()["diff"] + assert [x["ts_code"] for x in d["entered"]] == ["600418.SH"] + assert [x["ts_code"] for x in d["exited"]] == ["600519.SH"] + finally: + restore() + + +@case("[J7c] 上一版名册坏了要报故障, 不能显示成'首次'") +def t_j7c(): + import json as _j + from app.services import plan_feed as pf + store, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + pf.snapshot(plan(main=[_r("600418.SH", 1, 1.0)], date="2026-07-30")) + store["rows"][0]["roster_json"] = _j.dumps([]) # 名册没了, n_main 还写着 1 + out = pf.preview_changes(plan(main=[_r("600519.SH", 1, 2.0)], date="2026-07-31")) + assert out["diff"]["prev_broken"] is True + assert "读不出来" in out["hint"] + finally: + restore() + + +@case("[J7d] 持仓代码是前缀式也要认得出来") +def t_j7d(): + from app.services import plan_feed as pf + _, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + sys.modules["app.repo.pms_repo"].list_positions = \ + lambda only_open=False: [{"ts_code": "SH600418"}] # 前缀式 + pf.snapshot(plan(main=[_r("600418.SH", 1, 1.0)], date="2026-07-30")) + out = pf.preview_changes(plan(main=[_r("600519.SH", 1, 2.0)], date="2026-07-31")) + assert out["diff"]["held"]["n_watch"] == 1, "不归一的后果是持仓视图静默为空" + finally: + restore() + + +@case("[J8] 快照台账数得出同一计划日有几版") +def t_j8(): + from app.services import plan_feed as pf + _, restore = _with_fake_repo({"PMS_PLAN_SNAPSHOT": True}) + try: + pf.snapshot(plan(main=[_r("600418.SH", 1, 1.0)], date="2026-07-30")) + pf.snapshot(plan(main=[_r("600418.SH", 2, 1.0)], date="2026-07-30")) + pf.snapshot(plan(main=[_r("600418.SH", 1, 1.0)], date="2026-07-29")) + log = pf.snapshot_log() + assert log["ok"] and log["versions_per_date"]["2026-07-30"] == 2 + assert log["revised"] == {"2026-07-30": 2}, "同一 date 多版是 §7.3 那个现象, 要数出来" + finally: + restore() + + +def main(): + import logging + logging.disable(logging.CRITICAL) + passed, failed = 0, 0 + for name, fn in RESULTS: + try: + fn() + print(f" PASS {name}") + passed += 1 + except Exception: + print(f" FAIL {name}") + traceback.print_exc() + failed += 1 + print("-" * 60) + if failed: + print(f"FAILED: {failed} / {passed + failed}") + sys.exit(1) + print(f"ALL PASS ({passed} cases)") + + +if __name__ == "__main__": + main()