This commit is contained in:
zlt 2026-07-31 13:28:10 +08:00
parent a22875573d
commit 582925732d
16 changed files with 1951 additions and 29 deletions

89
Makefile Normal file
View File

@ -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

View File

@ -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 通道实现清单

View File

@ -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`
与持仓票视图。账本还空着, 所以持仓那段现在必然是空的, 等账本重建后才有内容。

334
app/core/plan_diff.py Normal file
View File

@ -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}"

View File

@ -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)` 加了唯一键当去重手段, 结果是: 同一天上游 ABA 地改回去时,
第三次的 A 因为"曾经存过"而落不下去, 库里最新仍是 B 于是页面把 BA 这次**真实的
回退**显示成 AB, 方向整个反了, 还会顺带报一句"手上这份没落进快照 (落库失败?)"
现在唯一键去掉, 只跟最新那行比: 变了就存, 没变就跳过表因此可能出现同 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

View File

@ -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

View File

@ -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。"""

View File

@ -576,6 +576,136 @@
{{ planCand.st_unknown.length }} 只无名称, ST 判不了已放行</span>
<br><span class="mono">{{ (planCand.items||[]).map(x=>x.ts_code).join(' ') }}</span>
</div>
<!-- 榜单变化 (PMS 自算; 上游 changes 恒为 null) -->
<el-divider content-position="left">榜单变化</el-divider>
<div class="row" style="margin-bottom:6px">
<el-button size="small" :loading="chgLoading" @click="loadChanges">重算变化</el-button>
<span class="muted" style="margin-left:8px" v-if="chgSt.hint">{{ chgSt.hint }}</span>
</div>
<el-alert v-if="chgSt.in_sync === false" type="error" effect="dark" show-icon
:closable="false" style="margin-bottom:8px"
title="手上这份计划没落进快照 —— 下面比的是两个旧版本, 与当前候选池不是同一份">
</el-alert>
<el-alert v-else-if="chg && chg.kind === 'same_date_revision'" type="warning" effect="dark"
show-icon :closable="false" style="margin-bottom:8px"
:title="'同一计划日 ' + chg.curr_date + ' 的新版本 —— 上游重算过, '
+ '此前拿到的候选池与现在不是同一份 (上游 /plan 是实时算的, 没有版本戳)'">
</el-alert>
<el-alert v-if="chg && chg.prev_broken" type="error" effect="dark" show-icon
:closable="false" style="margin-bottom:8px"
title="上一版名册读不出来 (快照损坏) —— 这次比不了, 下面的空白不代表没有变化">
</el-alert>
<el-alert v-if="chgHeld.n_watch" type="error" effect="dark" show-icon :closable="false"
style="margin-bottom:8px"
:title="'持仓票 ' + chgHeld.n_watch + ' 只掉榜 / 降档 / 丢了券商覆盖 —— 见下表红字'">
</el-alert>
<el-alert v-if="chg && chg.date_regressed" type="warning" effect="dark" show-icon
:closable="false" style="margin-bottom:8px"
title="计划日往回走了 (拉过一份旧计划?) —— 新进/掉榜的方向是反的">
</el-alert>
<div class="muted" v-if="chg" style="margin-bottom:8px">
{{ 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 }}
<span v-if="(chg.tail_churn.entered + chg.tail_churn.exited)">·
榜尾进出 {{ chg.tail_churn.entered + chg.tail_churn.exited }} 只未计入 (那一版吃满 top,
尾部是截断噪音不是观点变化)</span>
<span v-if="chg.tail_churn.held_exempt">· 其中
<b>{{ chg.tail_churn.held_exempt }} 只因是持仓票照报</b> (持仓票不受榜尾闸约束 ——
漏报一条就是一个仓位)</span>
<span v-if="chgRevised.length">· <b>同一计划日出现多版</b>:
{{ chgRevised.map(x => x[0] + '×' + x[1]).join(', ') }}</span>
</div>
<el-tabs v-model="chgTab" v-if="chg">
<el-tab-pane name="held" :label="'持仓票 (' + chgHeld.n_total + ')'">
<el-table :data="chgWatch.filter(x => x.held)" size="small" border max-height="220"
empty-text="持仓票没被摘也没降档">
<el-table-column prop="ts_code" label="代码" width="104"></el-table-column>
<el-table-column prop="name" label="名称" width="92"></el-table-column>
<el-table-column prop="why" label="发生了什么" min-width="200">
<template #default="s"><span style="color:#F56C6C">{{ s.row.why }}</span></template>
</el-table-column>
<el-table-column prop="rank" label="上版 rank" width="90"></el-table-column>
<el-table-column prop="theme" label="主题" width="120"></el-table-column>
</el-table>
<div class="muted" style="margin-top:6px">
上游把一只<b>持仓票</b>摘出榜或降了档, 是"还该不该继续拿着"的直接信号 ——
单独拎出来免得被几十条新进榜淹掉。<b>它只是提示, 不自动产生任何指令</b>:
要减要清仍走命令台或提议确认。
</div>
</el-tab-pane>
<el-tab-pane name="watch" :label="'掉榜与降档 (' + chgWatch.length + ')'">
<el-table :data="chgWatch" size="small" border max-height="220" empty-text="无">
<el-table-column prop="ts_code" label="代码" width="104"></el-table-column>
<el-table-column prop="name" label="名称" width="92"></el-table-column>
<el-table-column prop="why" label="发生了什么" min-width="180"></el-table-column>
<el-table-column prop="rank" label="上版 rank" width="90"></el-table-column>
<el-table-column label="持仓" width="70">
<template #default="s">
<el-tag v-if="s.row.held" type="danger" size="small">持仓</el-tag>
</template>
</el-table-column>
</el-table>
</el-tab-pane>
<el-tab-pane name="in" :label="'新进榜 (' + chg.counts.entered + ')'">
<el-table :data="chg.entered" size="small" border max-height="220" empty-text="无">
<el-table-column prop="rank" label="rank" width="62"></el-table-column>
<el-table-column prop="ts_code" label="代码" width="104"></el-table-column>
<el-table-column prop="name" label="名称" width="92"></el-table-column>
<el-table-column prop="tier" label="传导档" width="82"></el-table-column>
<el-table-column prop="theme" label="主题" width="120"></el-table-column>
<el-table-column prop="score" label="score" width="82"></el-table-column>
</el-table>
</el-tab-pane>
<el-tab-pane name="up" :label="'升档 (' + chg.counts.tier_up + ')'">
<el-table :data="chg.tier_up" size="small" border max-height="220" empty-text="无">
<el-table-column prop="rank" label="rank" width="62"></el-table-column>
<el-table-column prop="ts_code" label="代码" width="104"></el-table-column>
<el-table-column prop="name" label="名称" width="92"></el-table-column>
<el-table-column label="档位" min-width="150">
<template #default="s">{{ s.row.tier_from }} → <b>{{ s.row.tier_to }}</b></template>
</el-table-column>
<el-table-column prop="theme" label="主题" width="120"></el-table-column>
</el-table>
</el-tab-pane>
<el-tab-pane name="jump" :label="'名次跳变 (' + chg.counts.rank_jump + ')'">
<el-table :data="chg.rank_jump" size="small" border max-height="220" empty-text="无">
<el-table-column prop="ts_code" label="代码" width="104"></el-table-column>
<el-table-column prop="name" label="名称" width="92"></el-table-column>
<el-table-column label="名次" min-width="150">
<template #default="s">
{{ s.row.rank_from }} → {{ s.row.rank_to }}
<span :style="{color: s.row.rank_delta > 0 ? '#F56C6C' : '#67C23A'}">
({{ s.row.rank_delta > 0 ? '+' : '' }}{{ s.row.rank_delta }})</span>
</template>
</el-table-column>
<el-table-column prop="tier" label="传导档" width="82"></el-table-column>
<el-table-column prop="theme" label="主题" width="120"></el-table-column>
</el-table>
</el-tab-pane>
<el-tab-pane name="log" :label="'快照台账 (' + (chgLog.rows||[]).length + ')'">
<el-table :data="chgLog.rows" size="small" border max-height="220" empty-text="还没有快照">
<el-table-column prop="plan_date" label="计划日" width="104"></el-table-column>
<el-table-column prop="fetched_at" label="拉到的时刻" width="170"></el-table-column>
<el-table-column prop="n_main" label="主榜" width="70"></el-table-column>
<el-table-column prop="n_observe" label="观察" width="70"></el-table-column>
<el-table-column label="吃满 top" width="90">
<template #default="s">
<el-tag v-if="s.row.capped_main" type="warning" size="small">主榜</el-tag>
</template>
</el-table-column>
<el-table-column prop="digest" label="指纹" min-width="140"></el-table-column>
</el-table>
<div class="muted" style="margin-top:6px">
<b>同一「计划日」出现多行 = 上游那天重算过多版</b>(UPSTREAM_PLAN_API.md §7.3:
`/plan` 是实时算的, 应答里没有 generated_at, 下游本来没办法判断手上这份是哪一版)。
指纹不含 score —— 分数抖动但次序没变不算新版本, 否则表白涨而信息量为零。
</div>
</el-tab-pane>
</el-tabs>
<el-divider content-position="left">主榜明细</el-divider>
<el-table :data="planRows" size="small" border height="420">
<el-table-column prop="rank" label="rank" width="62"></el-table-column>
<el-table-column prop="ts_code" label="代码" width="104"></el-table-column>
@ -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');
</script>

View File

@ -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 列), 勿用。

View File

@ -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='上游计划名册快照 (榜单变化比对的底本)';

View File

@ -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 = [], []

View File

@ -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

View File

@ -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):

View File

@ -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():

View File

@ -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 全部单表合规")

View File

@ -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()