diff --git a/README.md b/README.md index 2be6f5e..a930019 100644 --- a/README.md +++ b/README.md @@ -60,7 +60,7 @@ scripts/ test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 11 例 test_batch5_units.py 决策系统信号流解析与消化口径 8 例 test_batch6_units.py ws 通道: 测试向量/签名/公钥/水位/DDL/逐笔入账 65 例 - test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 27 例 + test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 30 例 test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) 51 例 init_db.py 建表 (应用 ddl_pms_v1.sql, 幂等, 默认演练; 含 DDL 体检) check_db.py 实机连通性与表结构自检 (需真实 .env) @@ -199,7 +199,7 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ ## 已实现 / 待开发 -**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 232 例**。 +**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 235 例**。 ### 下一步(按可动工顺序) @@ -209,8 +209,8 @@ QMT ──trade/order_update──▶ pms-ws ──落 pms_qmt_inbox──▶ | 2 | **ws 通道联调**(协议 §9 的 S1/S2/S3) | 密钥已交换,可开工 | | 2.5 | 部署便利性:`make deploy` 一句话完成 build + up + 各 profile | 待办(每次改动要敲 4 条命令太烦) | | 2.6 | **行业硬拦截开闸**:跑一次「强刷 + 灌行业映射」把 `theme` 灌进 `pms_industry_map`,再把 `PMS_SECTOR_SOURCE` 改成 `custom_table` | 可做(数据源已就位,只差一次刷新 + 一个参数) | -| 2.7 | **主榜取全量**:上游默认只回 20 条(`counts.main=961`),候选池实际只在这 20 条里选 | 先跑 `probe_plan_api.py --try-limit` 探参数名(探到即填 `PMS_PLAN_QUERY_EXTRA`,无需改码);探不到则等上游给分页方式 | -| 2.8 | 上游计划的剩余待确认口径(6 条) | 阻塞:等上游答复,见 `UPSTREAM_PLAN_API.md` §4 | +| 2.7 | ~~主榜取全量~~ | ✅ 2026-07-30:拿到上游 `api.py`,`top`/`obs_top`/`theme_cap` 都是请求参数,已做成 PMS 三个显式参数;主题分散改由 `PMS_PLAN_THEME_CAP_LOCAL` 在候选阶段做 | +| 2.8 | 上游计划的剩余待确认口径(5 条) | 阻塞:等上游答复,见 `UPSTREAM_PLAN_API.md` §4 | | 3 | T0 做T(二期) | 可做,设计已有,无外部依赖 | | 4 | 择时实现 A(委托决策系统盘中择时) | 阻塞:等 bionic 侧接口 | | 5 | 研判闸接通 | 阻塞:等 bionic 侧 `process_intraday_audit` 新增 PMS 请求 direction。客户端已就位,接口好了在页面填 `PMS_JUDGE_API_BASE` 即通 | diff --git a/UPSTREAM_PLAN_API.md b/UPSTREAM_PLAN_API.md index ed4b2dc..4a84bd7 100644 --- a/UPSTREAM_PLAN_API.md +++ b/UPSTREAM_PLAN_API.md @@ -92,7 +92,13 @@ PMS_PLAN_MIN_SOURCES 0 evidence.n_sources 下限, 0=不设 PMS_PLAN_MIN_UPSIDE 0 预期空间下限 (0.5=+50%), 0=不设; 只过滤不排序 PMS_PLAN_STALE_TDAYS 1 日龄上限 (交易日) PMS_PLAN_THEME_SYNC true 刷新时把 theme 灌进 pms_industry_map -PMS_PLAN_QUERY_EXTRA (空) 附加查询串, 如 limit=1000 (取全量用, 见 §4 Q3) +PMS_PLAN_QUERY_EXTRA (空) 附加查询串逃生口 (上游加新参数时免改码); 同名键输给显式参数 + +发给上游的三个 (对方 api.py 的签名, 0=不传用上游默认): +PMS_PLAN_TOP 300 上游 top (上游默认 20) +PMS_PLAN_OBS_TOP 100 上游 obs_top (上游默认 10) +PMS_PLAN_THEME_CAP 999 上游 theme_cap (上游默认 5; 设大=让上游别裁) +PMS_PLAN_THEME_CAP_LOCAL 5 PMS 侧同主题限额, 在 TOP_N 截断**之前**生效, 0=不限 ``` 行业硬拦截的开法: 先跑一次「强刷 + 灌行业映射」(或等 08:40 调度位), 再把 @@ -177,7 +183,8 @@ sudo iptables -I INPUT -s 172.16.0.0/12 -p tcp --dport 8300 -j ACCEPT | Q5 | `counts.gate_covered` | md 里写作「全池**档位覆盖** 2386 只」——即被档位模型覆盖到的全池股票数, 与 main+observe(1068) 是包含关系。PMS 只透传展示 | | Q6a | `theme_cap=5` 的含义 | md 标题写明「主榜 Top 20(有券商预期、目标价不低于现价, **每主题限额 5**)」——上游**已自行执行**同主题限额。PMS 侧 `PMS_SECTOR_MAX_NAMES=4` 更严, 不冲突 | | — | 主榜的隐含前置 | 同一行标题揭示主榜已过两道筛: **有券商预期** + **目标价不低于现价**。所以主榜里不会出现 `upside < 0` | -| Q3a | 默认返回条数 | **主榜 20 / 观察档 10** (实测), 与 md 版「主榜 Top 20」一致。取全量的方式仍未知 —— 见下 Q3 | +| Q3 | 取全量的方式 | `top` / `obs_top` / `theme_cap` **全是请求参数** (上游 `api.py` 签名: `get_plan(date, format, top=20, obs_top=10, theme_cap=5)`)。20/10/5 是默认值不是政策 | +| Q6b | `theme_cap` 归谁管 | **归调用方**。既然是参数, 主题分散就该由 PMS 决定 —— 见 §5 | | — | 快照滞后是已知设计 | md 注: 「本日传导用的行情快照 = 2026-07-28(与计划日不同——**历史降级日口径**)」 | ### 仍需上游回答 @@ -186,25 +193,7 @@ sudo iptables -I INPUT -s 172.16.0.0/12 -p tcp --dport 8300 -j ACCEPT 而快照是 07-28。若 date 是生成日, 07-30 开盘该用哪一份? PMS 现在按「至多比今天旧 1 个 交易日」放行, 盘前 08:40 拉 —— 若上游出计划晚于 08:40, 这个调度位要往后挪。 -**Q3 (最重要) 怎么取主榜全量? 实测默认只回 20 条。** - -``` -counts(上游全量)={'main': 961, 'observe': 107} returned(本次应答)={'main': 20, 'observe': 10} -``` - -`format=md` 那版的标题也印证了: 「主榜 **Top 20**(有券商预期、目标价不低于现价, 每主题限额 5)」。 -两个后果得说清: - -1. **候选池实际只在这 20 条里选**, `PMS_PLAN_TOP_N=30` 根本吃不满 —— 排序池比以为的小 48 倍。 -2. **这 20 条已被上游按「每主题限额 5」裁过** (实测主题分布: 储能×5 / 传感器×5 / 整机×5 / - 整车×3 / 集成电路设计×2)。PMS 侧的行业集中度约束因此是在一个**已经被裁过**的池子上再裁 - 一次 —— 约束还成立, 但它看不到全貌, 也就选不出"上游主题限额之外但更合适"的票。 - -**请给取全量 (或分页) 的方式**: 参数名是什么? 有上限吗? 还是另有端点? -在此之前 PMS 侧已备好: 探参数名 `probe_plan_api.py --try-limit` (自动试 limit/top/size/ -page_size/… 十来个常见名), 探到就填 `PMS_PLAN_QUERY_EXTRA=limit=1000`, **不用改代码**。 -`parse_plan` 也已经把 `truncated` 标记算出来, 日志 warning + 页面抽屉橙色横幅都会显式提示, -不会静默拿 20 条当全量用。 +**~~Q3 怎么取主榜全量?~~ 已解决 (2026-07-30 拿到上游 `api.py`)。** 见下「§5 三个请求参数」。 **Q4 `changes` 的结构?** 样例是 `null`。若是「与上一份计划的差异」, 给个非空示例 —— PMS 想在页面上标「新进榜 / 掉榜」, 那是上游观点变化最直接的信号。 @@ -224,3 +213,59 @@ PMS 想在页面上标「新进榜 / 掉榜」, 那是上游观点变化最直 **Q10 (给决策系统侧) PMS 已不再读 `trading_buy_plan`。** 那张表还有别的消费方吗? 若没有, 上游可以停写 —— 少一处没有明确写入方的"事实源"。 + + +--- + +## 5. 三个请求参数与「谁来做主题分散」(2026-07-30) + +上游 `api.py` 的签名: + +```python +@app.get("/plan") +def get_plan(date: str | None = None, format: str = "json", + top: int = 20, obs_top: int = 10, theme_cap: int = 5): + data = plan.collect(date, top, obs_top, theme_cap) +``` + +**20 / 10 / 5 是默认值, 不是上游的既定政策。** 之前以为的「上游按每主题限额 5 裁过」, +其实是我们没传参数、吃了它的默认值。这一点改变了分工。 + +### 55 是怎么来的 + +`top=1000` 只回 55 条 —— 不是 top 卡的, 是 `theme_cap=5` 卡的: 过完「有券商预期 + +目标价不低于现价」之后大约 11 个主题, 每主题限 5 只 = 55。所以: + +``` +打分池 961 →(券商预期 + 目标价≥现价)→ ~11 个主题的若干只 →(theme_cap=5)→ 55 →(top)→ 实收 +``` + +`counts.main=961` 与实收条数**衡量的不是一回事**, 相减没有意义。判「还有没有更多」只能看 +**返回条数有没有吃满我们要的 top** —— 拿 961 去比会天天报假警。代码里 `_capped()` 就是这条。 + +### 分工: 上游别裁, PMS 自己裁 + +主题分散这件事放在哪一层做, 结果差很多: + +| 做法 | 后果 | +|---|---| +| 上游裁 (`theme_cap=5`) | 池子只有 55 只。PMS 看不到「上游主题限额之外但更合适」的票, 而且这个限额调不动 | +| 都不裁 | 池子宽了, 但 `PMS_PLAN_TOP_N=30` 是**按纯 score 切**的 —— 前 30 名可能全是储能, 切完交给规则闸, 规则闸按 `PMS_SECTOR_MAX_NAMES=4` 一拦就剩 4 只, **白瞎 26 个名额且日志上看不出来** | +| **上游别裁 + PMS 在候选阶段裁** ← 现在的做法 | 要个宽池子回来 (`theme_cap=999`), `select_candidates` 按 score 序走的时候就按 `PMS_PLAN_THEME_CAP_LOCAL` 摊开, **再**截 top_n。切出来的 30 只本身就是分散的, 不会被规则闸大批拦掉 | + +第三种在数学上不劣于第一种: 只要上游的每主题挑选也是 score 降序 (几乎必然), PMS 侧 +`THEME_CAP_LOCAL=5` 就能复现上游 `theme_cap=5` 的结果 —— 区别是这个数现在归我们调。 + +默认值: + +``` +PMS_PLAN_TOP=300 PMS_PLAN_OBS_TOP=100 PMS_PLAN_THEME_CAP=999 # 发给上游: 要宽池子 +PMS_PLAN_THEME_CAP_LOCAL=5 # PMS 侧摊开, 与原行为等价 +PMS_PLAN_TOP_N=30 # 摊开之后再截断 +``` + +`dropped.theme` 会记下被主题限额挡掉几只, 页面抽屉和探活脚本都会打「入池主题分布」—— +限额有没有真在起作用, 一眼能看出来。 + +> 注意: 服务器上如果还留着 `PMS_PLAN_QUERY_EXTRA=top=30`, 清掉它。同名键以显式参数为准, +> 留着不影响功能, 但两个地方写同一件事迟早看走眼。 diff --git a/app/services/param_store.py b/app/services/param_store.py index 6035ecd..8673b57 100644 --- a/app/services/param_store.py +++ b/app/services/param_store.py @@ -66,6 +66,10 @@ DESC = { "PMS_PLAN_API_PATH": "计划接口路径 (默认 /plan)", "PMS_PLAN_TIMEOUT": "计划接口超时 (秒)", "PMS_PLAN_CACHE_SEC": "计划缓存秒数 (上游日频产出)", + "PMS_PLAN_TOP": "发给上游的 top: 主榜要多少条 (上游默认 20; 0=不传)", + "PMS_PLAN_OBS_TOP": "发给上游的 obs_top: 观察档要多少条 (上游默认 10; 0=不传)", + "PMS_PLAN_THEME_CAP": "发给上游的 theme_cap: 上游侧同主题限额 (上游默认 5; 设大值=让上游别裁, 由 PMS 自己裁; 0=不传)", + "PMS_PLAN_THEME_CAP_LOCAL": "PMS 侧同主题限额, 在 TOP_N 截断之前按 score 序生效, 0=不限。防止宽池子里前 N 名被单一主题垄断、切完再被规则闸拦掉", "PMS_PLAN_TOP_N": "主榜按 score 降序取前 N 只进候选池", "PMS_PLAN_TIERS": "传导档白名单 (如 强传导), 留空=不按档过滤", "PMS_PLAN_INCLUDE_OBSERVE": "观察档是否进候选池", @@ -74,7 +78,7 @@ DESC = { "PMS_PLAN_MIN_UPSIDE": "候选预期空间下限 (相对现价, 0.5=+50%), 0=不设。只过滤不参与排序; 设了就会把 upside 缺失的行(含整个观察档)一起挡掉", "PMS_PLAN_STALE_TDAYS": "计划日龄超此交易日数即判过期拒用 (防上游停更时拿旧榜当今天)", "PMS_PLAN_THEME_SYNC": "刷新计划时把 evidence.theme 灌进 pms_industry_map (行业源 custom_table 的数据来源)", - "PMS_PLAN_QUERY_EXTRA": "计划接口附加查询串, 如 limit=1000。上游默认只回主榜 20 条(counts 却是 961), 用 probe_plan_api.py --try-limit 探出参数名后填这里", + "PMS_PLAN_QUERY_EXTRA": "计划接口附加查询串逃生口 (上游加了新参数时不用改代码); top/obs_top/theme_cap 请用各自的显式参数, 同名键以显式参数为准", "PMS_SECTOR_SOURCE": "行业划分数据源: 空=约束停用 / custom_table (推荐, 由上游计划的 theme 灌数) / gp_stock_category", "PMS_SECTOR_MAX_NAMES": "同行业最大持仓只数 (硬拦截)", "PMS_SECTOR_MAX_RATIO": "同行业最大占总仓比例 (硬拦截)", @@ -264,7 +268,9 @@ _RANGES = { "PMS_EXEC_WINDOW_TDAYS": (1, 20), "PMS_BRAKE_DAYS": (0, 30), "PMS_PLAN_TOP_N": (1, 1000), "PMS_PLAN_TIMEOUT": (1, 120), "PMS_PLAN_CACHE_SEC": (5, 86400), "PMS_PLAN_STALE_TDAYS": (0, 20), - "PMS_PLAN_MIN_SOURCES": (0, 100), + "PMS_PLAN_MIN_SOURCES": (0, 100), "PMS_PLAN_TOP": (0, 5000), + "PMS_PLAN_OBS_TOP": (0, 5000), "PMS_PLAN_THEME_CAP": (0, 5000), + "PMS_PLAN_THEME_CAP_LOCAL": (0, 1000), } diff --git a/app/services/plan_feed.py b/app/services/plan_feed.py index d4847e1..26d946b 100644 --- a/app/services/plan_feed.py +++ b/app/services/plan_feed.py @@ -50,6 +50,12 @@ logger = logging.getLogger("pms.plan") BUCKET_MAIN, BUCKET_OBSERVE = "main", "observe" +# 上游 GET /plan 的签名 (2026-07-30 拿到对方 api.py 确认): +# get_plan(date=None, format="json", top=20, obs_top=10, theme_cap=5) +# 三个都是**请求参数**, 不是上游的既定政策 —— 主榜给多少、观察档给多少、每主题限几只, +# 全由调用方 (也就是 PMS) 决定。默认值写在这里, 用来判断"是不是被条数卡住了"。 +UPSTREAM_DEFAULT_TOP, UPSTREAM_DEFAULT_OBS_TOP, UPSTREAM_DEFAULT_THEME_CAP = 20, 10, 5 + # 候选池来源 (PMS_CANDIDATE_SOURCE) SRC_PLAN_API, SRC_BUY_PLAN, SRC_BOTH = "plan_api", "buy_plan", "both" @@ -114,8 +120,12 @@ def _rows(raw, bucket: str) -> list: return out -def parse_plan(payload) -> dict: - """应答 → 内部结构。缺 date 或两档全空都算变形 (抛 PlanFeedError)。""" +def parse_plan(payload, *, requested=None) -> dict: + """应答 → 内部结构。缺 date 或两档全空都算变形 (抛 PlanFeedError)。 + + requested: 本次实际发出的 {top, obs_top} (缺省按上游签名的默认值 20/10 算)。 + 判"还有没有更多"必须拿它跟返回条数比, **不能拿 counts 比** —— 见 _capped 的注释。 + """ if not isinstance(payload, dict): raise PlanFeedError(f"应答不是 JSON 对象: {type(payload).__name__}") date = _text_or_none(payload.get("date")) @@ -126,6 +136,11 @@ def parse_plan(payload) -> dict: if not main and not observe: raise PlanFeedError(f"计划 {date} 主榜与观察档都是空的") counts = payload.get("counts") if isinstance(payload.get("counts"), dict) else {} + req = dict(requested or {}) + req_top = _int_or_none(req.get("top")) + req_top = UPSTREAM_DEFAULT_TOP if req_top is None else req_top + req_obs = _int_or_none(req.get("obs_top")) + req_obs = UPSTREAM_DEFAULT_OBS_TOP if req_obs is None else req_obs themes = {} for r in main + observe: # 主榜在前, 同码以主榜的 theme 为准 if r["theme"] and r["ts_code"] not in themes: @@ -140,18 +155,30 @@ def parse_plan(payload) -> dict: "observe": _int_or_none(counts.get("observe")), "gate_covered": _int_or_none(counts.get("gate_covered"))}, "returned": {"main": len(main), "observe": len(observe)}, - # counts 是上游全量, returned 是这次真给了几条 —— 两者不等就是被截断了。 - # 2026-07-30 实测: counts.main=961 而只回 20 条 (上游默认 Top20), 排序池远小于 - # 以为的规模, 且那 20 条已经被上游按"每主题限额 5"裁过。必须显式暴露, 不能当没事。 - "truncated": {"main": _truncated(counts.get("main"), len(main)), - "observe": _truncated(counts.get("observe"), len(observe))}, + # 漏斗: counts 是上游的**打分池规模**, returned 是过完 + # 「有券商预期 + 目标价不低于现价 + 每主题限额 + top」之后真给了几条。 + # 两个数衡量的不是一回事, 相减没有意义 —— 只做展示。 + "funnel": {"scored_main": _int_or_none(counts.get("main")), "returned_main": len(main), + "scored_observe": _int_or_none(counts.get("observe")), + "returned_observe": len(observe)}, + "requested": {"top": req_top, "obs_top": req_obs, + "theme_cap": _int_or_none(req.get("theme_cap"))}, + "truncated": {"main": _capped(len(main), req_top), + "observe": _capped(len(observe), req_obs)}, "main": main, "observe": observe, "themes": themes, } -def _truncated(total, got) -> bool: - t = _int_or_none(total) - return bool(t is not None and got < t) +def _capped(got: int, requested) -> bool: + """"是不是还有更多没拿到" = 返回条数吃满了我们要的条数。 + + **不能拿 counts 判。** counts.main=961 是打分池规模, 而 returned 是过完券商预期、 + 目标价不低于现价、每主题限额、top 之后的结果 —— 实测 top=1000 也只回 55 条 (被 + theme_cap=5 卡住)。拿 961 跟 55 比会永远报"被截断", 变成一个天天喊狼来了的假警报。 + 吃满才说明是条数卡的, 没吃满就是上游确实只有这么多能给。 + """ + r = _int_or_none(requested) + return bool(r is not None and r > 0 and got >= r) # ================================================================ 纯逻辑: 新鲜度 @@ -174,13 +201,19 @@ def assert_fresh(plan: dict, *, max_stale_tdays: int = 1, today=None) -> int: # ================================================================ 纯逻辑: 筛选 def select_candidates(plan: dict, *, held=(), black=(), top_n: int = 30, tiers=None, include_observe: bool = False, min_score=None, - min_sources: int = 0, min_upside=None) -> dict: + min_sources: int = 0, min_upside=None, theme_cap: int = 0) -> dict: """排序池 → 候选清单。 排序: score 降序, 同分按 rank 升序 (上游 rank 已是它自己的最终次序, 拿来当稳定次序)。 tier 白名单**只对带 tier 的行生效** —— 观察档没有 tier, 它的闸门是 include_observe。 min_upside 相反, **对所有行生效**: upside 缺失按 0 算一起挡掉 (观察档 upside 恒为 null, 所以设了下限等于把观察档全挡了)。方向选保守那边 —— 宁可少票。 + + theme_cap: 同主题最多取几只, **在 top_n 截断之前**按 score 序生效 (0=不限)。 + 这一层存在的理由: 上游的 theme_cap 是请求参数, 我们可以让它别裁 (要个宽池子), 但 + top_n 那一刀是按纯 score 切的 —— 宽池子里前 30 名可能全是储能, 切完再交给规则闸, + 规则闸按 PMS_SECTOR_MAX_NAMES 一拦就剩 4 只, 白瞎 26 个名额且**日志上看不出来**。 + 在候选阶段先按主题摊开, top_n 切出来的才是能用的票。 """ held = {normalize_code(c) for c in (held or []) if c} black = {normalize_code(c) for c in (black or []) if c} @@ -193,9 +226,10 @@ def select_candidates(plan: dict, *, held=(), black=(), top_n: int = 30, tiers=N if include_observe: pool += list(plan.get("observe") or []) + theme_cap = int(theme_cap or 0) dropped = {"held": 0, "black": 0, "tier": 0, "score": 0, "sources": 0, "upside": 0, - "dup": 0, "capped": 0} - passed, seen = [], set() + "theme": 0, "dup": 0, "capped": 0} + passed, seen, per_theme = [], set(), {} for r in sorted(pool, key=lambda x: (-(x.get("score") or 0.0), x.get("rank") or 10 ** 9)): c = r["ts_code"] if c in seen: @@ -220,6 +254,12 @@ def select_candidates(plan: dict, *, held=(), black=(), top_n: int = 30, tiers=N if min_upside is not None and (r.get("upside") or 0.0) < min_upside: dropped["upside"] += 1 continue + if theme_cap > 0: + t = r.get("theme") or "(无主题)" + if per_theme.get(t, 0) >= theme_cap: + dropped["theme"] += 1 + continue + per_theme[t] = per_theme.get(t, 0) + 1 passed.append(r) n = max(0, int(top_n or 0)) or len(passed) @@ -253,11 +293,35 @@ def _params() -> dict: "min_upside": ps.get_float("PMS_PLAN_MIN_UPSIDE", 0.0), "stale_tdays": ps.get_int("PMS_PLAN_STALE_TDAYS", 1), "theme_sync": ps.get_bool("PMS_PLAN_THEME_SYNC", True), + "top": ps.get_int("PMS_PLAN_TOP", 300), + "obs_top": ps.get_int("PMS_PLAN_OBS_TOP", 100), + "theme_cap": ps.get_int("PMS_PLAN_THEME_CAP", 999), + "theme_cap_local": ps.get_int("PMS_PLAN_THEME_CAP_LOCAL", 5), "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(), } +def build_query(p: dict) -> dict: + """发给上游的查询参数。显式参数覆盖 PMS_PLAN_QUERY_EXTRA 里的同名键。 + + QUERY_EXTRA 保留是为了上游哪天加了新参数时不用改代码; 但 top/obs_top/theme_cap 这三个 + 已经知道签名了, 走各自的显式参数 —— 同一件事有两个入口的时候, 得有个明确的赢家。 + """ + q = dict(p.get("query_extra") or {}) + for key, name in (("top", "top"), ("obs_top", "obs_top"), ("theme_cap", "theme_cap")): + v = int(p.get(key) or 0) + if v > 0: + q[name] = str(v) # 0 = 不传该参数, 用上游默认 + return q + + +def effective_query() -> dict: + """当前参数下实际会发出去的查询串。探活脚本要跟生产走同一条路 —— 上一版就是因为 + 没带这三个参数, 报出来的"请求参数"是上游默认值 20/10/5, 跟真实抓取行为对不上。""" + return build_query(_params()) + + def enabled() -> bool: return bool(_params()["base"]) @@ -313,7 +377,7 @@ def fetch(*, date=None, base=None, path=None, timeout=None, extra_params=None) - raise except Exception as e: raise PlanFeedError(f"拉取上游计划失败 {url}: {type(e).__name__}: {e}") from e - plan = parse_plan(payload) + plan = parse_plan(payload, requested=q) plan["url"] = url plan["fetched_at"] = time.time() plan["requested_date"] = date @@ -323,7 +387,8 @@ def fetch(*, date=None, base=None, path=None, timeout=None, extra_params=None) - def get_plan(*, force: bool = False, date=None) -> dict: """带缓存的当前计划。失败同样缓存 FAIL_CACHE_SEC, 但每次调用都照样抛。""" p = _params() - key = (p["base"], p["path"], date or "", tuple(sorted(p["query_extra"].items()))) + q = build_query(p) + key = (p["base"], p["path"], date or "", tuple(sorted(q.items()))) now = time.time() with _lock: fresh_hit = (not force and _cache["key"] == key and _cache["plan"] is not None @@ -334,7 +399,7 @@ def get_plan(*, force: bool = False, date=None) -> dict: and now - _cache["at"] < FAIL_CACHE_SEC): raise PlanFeedError(_cache["error"]) try: - plan = fetch(date=date, extra_params=p["query_extra"]) + plan = fetch(date=date, extra_params=q) assert_fresh(plan, max_stale_tdays=p["stale_tdays"]) except PlanFeedError as e: with _lock: @@ -349,11 +414,9 @@ def get_plan(*, force: bool = False, date=None) -> dict: plan["date"], plan["returned"]["main"], plan["returned"]["observe"], plan["age_tdays"], plan["url"]) if plan["truncated"]["main"]: - logger.warning("[上游计划] 主榜被截断: 上游 counts=%s 但只回了 %d 条 —— 候选池实际是从" - " 这 %d 条里选, 且上游已按自己的主题限额裁过。取全量的参数名探出来后" - " 填进 PMS_PLAN_QUERY_EXTRA (如 limit=1000)", - plan["counts"]["main"], plan["returned"]["main"], - plan["returned"]["main"]) + logger.warning("[上游计划] 主榜正好吃满 top=%s (打分池 %s) —— 可能还有更多没拿到, " + "调大 PMS_PLAN_TOP 再看", plan["requested"]["top"], + plan["funnel"]["scored_main"]) return plan @@ -393,7 +456,8 @@ def candidates(*, held=(), black=()) -> dict: include_observe=p["include_observe"], min_score=(p["min_score"] or None), min_sources=p["min_sources"], - min_upside=(p["min_upside"] or None)) + min_upside=(p["min_upside"] or None), + theme_cap=p["theme_cap_local"]) out["age_tdays"] = plan.get("age_tdays") out["theme_cap"] = plan.get("theme_cap") return out @@ -404,6 +468,7 @@ def status() -> dict: p = _params() st = {"source": p["source"], "base": p["base"], "path": p["path"], "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"])} if not p["base"]: @@ -422,6 +487,7 @@ def status() -> dict: "heat_date": plan.get("heat_date"), "market_snapshot_days": plan.get("market_snapshot_days"), "counts": plan["counts"], "returned": plan["returned"], + "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"), diff --git a/app/web/static/index.html b/app/web/static/index.html index 0f8cbaa..a1a366d 100644 --- a/app/web/static/index.html +++ b/app/web/static/index.html @@ -556,16 +556,16 @@
{{ planResult }}
diff --git a/config/settings.py b/config/settings.py
index a050388..01d9bc0 100644
--- a/config/settings.py
+++ b/config/settings.py
@@ -82,6 +82,13 @@ class Settings(BaseSettings):
PMS_PLAN_API_PATH: str = "/plan"
PMS_PLAN_TIMEOUT: int = 10 # 单次请求超时 (秒)
PMS_PLAN_CACHE_SEC: int = 300 # 计划缓存秒数 (上游日频产出, 没必要每跳都拉)
+ # 下面三个是**发给上游的请求参数** (对方 api.py: top=20/obs_top=10/theme_cap=5 都是
+ # 可传的, 不是它的既定政策)。0 = 不传, 用上游默认。要个宽池子回来, 主题分散交给
+ # PMS_PLAN_THEME_CAP_LOCAL 在候选阶段做 —— 理由见 plan_feed.select_candidates。
+ PMS_PLAN_TOP: int = 300 # 上游 top: 主榜要多少条
+ PMS_PLAN_OBS_TOP: int = 100 # 上游 obs_top: 观察档要多少条
+ PMS_PLAN_THEME_CAP: int = 999 # 上游 theme_cap: 让它别裁 (PMS 自己裁)
+ PMS_PLAN_THEME_CAP_LOCAL: int = 5 # PMS 侧同主题限额, 在 TOP_N 截断**之前**生效, 0=不限
PMS_PLAN_TOP_N: int = 30 # 主榜按 score 降序取前 N 进池 (900+ 只是排序池不是清单)
PMS_PLAN_TIERS: str = "强传导" # 传导档白名单; 留空=不按档过滤
PMS_PLAN_INCLUDE_OBSERVE: bool = False # 观察档是否进候选池 (上游把它定位为备选, 无 upside)
@@ -91,8 +98,8 @@ class Settings(BaseSettings):
# 口径, 噪音大 (榜首能到 +212%), 只做下限过滤, **不参与排序** —— 排序始终是 score
PMS_PLAN_STALE_TDAYS: int = 1 # 计划日龄超此交易日数即判过期并拒用 (防上游停更)
PMS_PLAN_THEME_SYNC: bool = True # 刷新时把 evidence.theme 灌进 pms_industry_map
- PMS_PLAN_QUERY_EXTRA: str = "" # 附加查询串, 如 "limit=1000"。上游默认只回主榜
- # 20 条 (counts 却是 961), 取全量的参数名待确认 —— 探出来填这里即可, 不用改代码
+ PMS_PLAN_QUERY_EXTRA: str = "" # 附加查询串逃生口, 如 "foo=1"。上游哪天加了新参数
+ # 不用改代码即可透传; top/obs_top/theme_cap 已有显式参数, 同名键以显式参数为准
# --- 行业约束 (硬拦截; 数据源接口化) ---
PMS_SECTOR_SOURCE: str = "" # "" = 停用并页面提示 / custom_table / gp_stock_category
diff --git a/scripts/probe_plan_api.py b/scripts/probe_plan_api.py
index 6179b41..598f473 100644
--- a/scripts/probe_plan_api.py
+++ b/scripts/probe_plan_api.py
@@ -182,21 +182,27 @@ def main():
print(" → PMS_PLAN_API_BASE 为空。页面「参数设置」填上再跑。")
return 2
+ q = pf.effective_query()
+ print(f" 发出的查询参数 {q or '(无 —— 将吃上游默认 top=20/obs_top=10/theme_cap=5)'}")
try:
- plan = pf.fetch(date=args.date, base=base)
+ plan = pf.fetch(date=args.date, base=base, extra_params=q)
except pf.PlanFeedError as e:
print(f"[2] 取数失败: {e}")
return 1
print(f"[2] 取数成功 {plan['url']}")
print(f" date={plan['date']} heat_date={plan['heat_date']} "
f"snapshot={plan['market_snapshot_days']} theme_cap={plan['theme_cap']}")
- print(f" counts(上游全量)={plan['counts']} returned(本次应答)={plan['returned']}")
+ print(f" 请求参数 {plan['requested']}")
+ f = plan["funnel"]
+ print(f" 漏斗: 打分池 主榜 {f['scored_main']} / 观察 {f['scored_observe']}"
+ f" →(券商预期 + 目标价≥现价 + 主题限额 + top)→"
+ f" 实收 主榜 {f['returned_main']} / 观察 {f['returned_observe']}")
if plan["truncated"]["main"]:
- print(f" ! 主榜被截断: counts={plan['counts']['main']} 但只回了 "
- f"{plan['returned']['main']} 条")
- print(f" 候选池实际只在这 {plan['returned']['main']} 条里选, 且它们已被上游按"
- f"「每主题限额 {plan['theme_cap']}」裁过。")
- print(f" 探取全量的参数名: probe_plan_api.py --try-limit")
+ print(f" ! 主榜正好吃满 top={plan['requested']['top']} —— 可能还有更多, "
+ f"调大 PMS_PLAN_TOP 再看")
+ else:
+ print(f" 主榜没吃满 top={plan['requested']['top']}, 说明这就是上游能给的全部"
+ f" (受主题限额 {plan['theme_cap']} 与价格筛限制)")
age = pf.plan_age_tdays(plan["date"])
limit = ps.get_int("PMS_PLAN_STALE_TDAYS", 1)
@@ -238,12 +244,20 @@ def main():
plan, top_n=ps.get_int("PMS_PLAN_TOP_N", 30), tiers=ps.get_list("PMS_PLAN_TIERS", []),
include_observe=ps.get_bool("PMS_PLAN_INCLUDE_OBSERVE", False),
min_score=(ps.get_float("PMS_PLAN_MIN_SCORE", 0.0) or None),
- min_sources=ps.get_int("PMS_PLAN_MIN_SOURCES", 0))
+ min_sources=ps.get_int("PMS_PLAN_MIN_SOURCES", 0),
+ min_upside=(ps.get_float("PMS_PLAN_MIN_UPSIDE", 0.0) or None),
+ theme_cap=ps.get_int("PMS_PLAN_THEME_CAP_LOCAL", 5))
print(f"[6] 按当前参数筛选 (top_n={ps.get_int('PMS_PLAN_TOP_N', 30)} "
f"tiers={ps.get_list('PMS_PLAN_TIERS', [])} "
- f"observe={ps.get_bool('PMS_PLAN_INCLUDE_OBSERVE', False)})")
+ f"observe={ps.get_bool('PMS_PLAN_INCLUDE_OBSERVE', False)} "
+ f"PMS侧主题限额={ps.get_int('PMS_PLAN_THEME_CAP_LOCAL', 5)})")
print(f" 排序池 {sel['considered']} → 合格 {sel['eligible']} → 取 {len(sel['items'])} 只")
print(f" 丢弃明细 {sel['dropped']} (此处未扣持仓/黑名单, 下命令时还会再扣)")
+ st = {}
+ for x in sel["items"]:
+ st[x.get("theme") or "(无)"] = st.get(x.get("theme") or "(无)", 0) + 1
+ print(" 入池主题分布: " + ", ".join(f"{k}×{v}" for k, v in
+ sorted(st.items(), key=lambda y: -y[1])))
print(" " + ", ".join(x["ts_code"] for x in sel["items"]))
if args.json:
diff --git a/scripts/run_tests.py b/scripts/run_tests.py
index 2633db3..99755a1 100644
--- a/scripts/run_tests.py
+++ b/scripts/run_tests.py
@@ -11,9 +11,9 @@
test_batch4_units.py 动作引擎 四类自主动作触发与数量口径 (11 例)
test_batch5_units.py 决策系统信号流解析与消化口径 (8 例)
test_batch6_units.py ws 通道: 测试向量/签名/公钥/水位/DDL/逐笔入账 (65 例)
- test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 (27 例)
+ test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 (30 例)
test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (51 例)
- 共 232 例
+ 共 235 例
任一子集失败即整体失败 (退出码 1)。
"""
import os
diff --git a/scripts/test_batch7_units.py b/scripts/test_batch7_units.py
index 12f1712..225fb9f 100644
--- a/scripts/test_batch7_units.py
+++ b/scripts/test_batch7_units.py
@@ -155,16 +155,30 @@ def _():
assert r["n_sources"] == 7 and r["theme"] == "整车" # theme 两头空白要去掉
-@case("解析·截断标记: counts 大于实际返回条数就是被上游截断了 (2026-07-30 实测 961→20)")
+@case("解析·漏斗: counts 是打分池, returned 是过筛后 —— 两者相减没有意义, 只做展示")
def _():
- p = pf.parse_plan(SAMPLE) # counts.main=961, 实际 6 条
- assert p["truncated"] == {"main": True, "observe": True}, p["truncated"]
- full = {"date": "2026-07-29", "counts": {"main": 6, "observe": 3},
- "main": SAMPLE["main"], "observe": SAMPLE["observe"]}
- assert pf.parse_plan(full)["truncated"] == {"main": False, "observe": False}
- # counts 缺失时不能瞎报截断 (无从判断)
- nc = {"date": "2026-07-29", "main": SAMPLE["main"]}
- assert pf.parse_plan(nc)["truncated"] == {"main": False, "observe": False}
+ p = pf.parse_plan(SAMPLE)
+ assert p["funnel"] == {"scored_main": 961, "returned_main": 6,
+ "scored_observe": 107, "returned_observe": 3}, p["funnel"]
+
+
+@case("解析·吃满才叫截断: 拿 counts 判会永远报警 (961 vs 55 是两把尺子)")
+def _():
+ # 上游签名默认 top=20 / obs_top=10。样例只有 6+3 条 —— 没吃满, 就是上游能给的全部,
+ # 尽管 counts 写着 961。**这正是不能拿 counts 判截断的原因。**
+ p = pf.parse_plan(SAMPLE)
+ assert p["requested"] == {"top": 20, "obs_top": 10, "theme_cap": None}, p["requested"]
+ assert p["truncated"] == {"main": False, "observe": False}, p["truncated"]
+ # 要了 6 条正好给 6 条 → 吃满, 可能还有更多
+ q = pf.parse_plan(SAMPLE, requested={"top": 6, "obs_top": 3, "theme_cap": 999})
+ assert q["truncated"] == {"main": True, "observe": True}, q["truncated"]
+ assert q["requested"]["theme_cap"] == 999
+ # 要 100 条只给 6 条 → 没吃满
+ r = pf.parse_plan(SAMPLE, requested={"top": 100, "obs_top": 100})
+ assert r["truncated"] == {"main": False, "observe": False}
+ # top=0 (不传) 时无从判断, 不许瞎报
+ z = pf.parse_plan(SAMPLE, requested={"top": 0, "obs_top": 0})
+ assert z["truncated"] == {"main": False, "observe": False}
@case("取数·附加查询串: 解析 / 合并进请求 / date 不可被覆盖 / 坏串忽略不炸")
@@ -309,6 +323,39 @@ def _():
assert [x["ts_code"] for x in r["items"]] == ["600001.SH", "600002.SH"], r["items"]
+@case("筛选·候选级主题限额在 top_n 截断之前生效 (否则前 N 名被单一主题垄断)")
+def _():
+ # 造一个"宽池子": 储能 6 只分最高, 传感器 3 只, 整车 2 只
+ rows = ([_m(i, f"SH60{i:04d}", f"储{i}", 300 - i, "储能", 0.2, 1.0) for i in range(1, 7)]
+ + [_m(10 + i, f"SH61{i:04d}", f"传{i}", 200 - i, "传感器", 0.2, 1.0) for i in range(1, 4)]
+ + [_m(20 + i, f"SH62{i:04d}", f"整{i}", 100 - i, "整车", 0.2, 1.0) for i in range(1, 3)])
+ p = pf.parse_plan({"date": "2026-07-29", "main": rows})
+ # 不限主题: 前 5 名全是储能 —— 规则闸按 PMS_SECTOR_MAX_NAMES 一拦就废掉大半
+ plain = pf.select_candidates(p, top_n=5)
+ assert {x["theme"] for x in plain["items"]} == {"储能"}, plain["items"]
+ # 限 2 只/主题: 先摊开再截断, 5 个名额分到三个主题
+ capped = pf.select_candidates(p, top_n=5, theme_cap=2)
+ got = [(x["theme"], x["ts_code"]) for x in capped["items"]]
+ assert [t for t, _ in got] == ["储能", "储能", "传感器", "传感器", "整车"], got
+ assert capped["dropped"]["theme"] == 5 # 储能多 4 只 + 传感器多 1 只
+ # 主题为空的行归到「(无主题)」一档, 不跟着别的主题挤
+ p2 = pf.parse_plan({"date": "2026-07-29", "main": [
+ {"rank": 1, "code": "SH600001", "score": 9}, {"rank": 2, "code": "SH600002", "score": 8}]})
+ assert len(pf.select_candidates(p2, top_n=5, theme_cap=1)["items"]) == 1
+
+
+@case("取数·build_query: 显式 top/obs_top/theme_cap 覆盖 QUERY_EXTRA 同名键; 0=不传")
+def _():
+ assert pf.build_query({"top": 300, "obs_top": 100, "theme_cap": 999, "query_extra": {}}) == \
+ {"top": "300", "obs_top": "100", "theme_cap": "999"}
+ # 0 = 不传该参数 (用上游默认)
+ assert pf.build_query({"top": 0, "obs_top": 0, "theme_cap": 0, "query_extra": {}}) == {}
+ # QUERY_EXTRA 里的同名键输给显式参数; 非同名键照旧透传
+ q = pf.build_query({"top": 300, "obs_top": 0, "theme_cap": 0,
+ "query_extra": {"top": "20", "foo": "bar"}})
+ assert q == {"top": "300", "foo": "bar"}, q
+
+
@case("筛选·输出字段: sector 用 theme 灌 (planner 吃这个), score 缺失兜 0.0 不留 None")
def _():
d = {"date": "2026-07-29", "main": [{"rank": 1, "code": "SH600418",
@@ -326,11 +373,11 @@ def _():
_m(500, "SH600519", "弱票", 220.0, "白酒", 0.1, 0.5, tier="弱传导")]
p = pf.parse_plan(d)
r = pf.select_candidates(p, held=["SH600418"], black=["SZ300952"], tiers=["强传导"],
- min_sources=7, top_n=2, include_observe=True)
+ min_sources=7, top_n=2, include_observe=True, theme_cap=2)
dr = r["dropped"]
assert r["considered"] == 10, r["considered"]
assert r["considered"] == r["eligible"] + dr["held"] + dr["black"] + dr["tier"] \
- + dr["score"] + dr["sources"] + dr["upside"] + dr["dup"], (r, dr)
+ + dr["score"] + dr["sources"] + dr["upside"] + dr["theme"] + dr["dup"], (r, dr)
assert len(r["items"]) == r["eligible"] - dr["capped"] == 2