diff --git a/README.md b/README.md index 649d7b0..6236e69 100644 --- a/README.md +++ b/README.md @@ -298,7 +298,34 @@ make t-gate # 随时: 规则闸/研判闸拒了什么、为什么 ## 已实现 / 待开发 -**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**榜单变化提示**(名册快照 + 新进/掉榜/档位升降/覆盖翻转/名次跳变,持仓票单列,榜尾截断噪音闸);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 351 例**。 +**已实现**:建表 DDL 与建表脚本;配置与运行参数中心;仓位规划器与安全垫账;命令系统(27 类命令全目录 + 双状态机 + 冲突识别);方案生成器(降仓凑额四档、升仓、建仓分批、清仓/减至、行业清仓与限额、暂停买入撤单);账本回放与对账引擎(成交认领、外部成交并入 BASE 告警、以下游为准修正、除权检测、T+1 可用量、连续不一致升级);规则闸终检;择时执行器实现 B(分日配额、分笔、VWAP/回踩/不追高、14:45 兜底、停牌一字板顺延、窗口耗尽收口、挂单有效期);动作引擎四类自主动作 + 研判闸客户端 + 提议分流;决策系统信号消化(两条流独立消费组订阅、置信度分档转清仓指令或提议);管理页面四块 + 运维/日报抽屉;调度器九个调度位;**上游选股计划接口接入**(`/plan` 取候选池、交易日龄硬校验、`theme` 灌行业映射表、页面预览抽屉与不可用横幅);**榜单变化提示**(名册快照 + 新进/掉榜/档位升降/覆盖翻转/名次跳变,持仓票单列,榜尾截断噪音闸);**ws 直连通道的连接层**(常驻进程 + 出口队列 + 签名 + seq 水位与累积确认,见下);**单测 378 例**。 + +### 静默失败专项(2026-07-31,八条已修) + +联测走到「记账并下发」这一步时,接连撞上几处「命令说完成了、实际什么都没发生」,于是停下来专门查了一轮。查出来的**不是八个孤立的 bug,是同一种病**: + +> 本系统里大量写入函数**失败时不抛异常,只回 `{"ok": False, ...}` 或影响 0 行**。返回值一丢,写入没发生,而调用方照常往下走、照常回 `ok=True`、页面照常显示「已完成」。**失败长得像成功。** + +八条实例,按「一旦发生会怎样」排序: + +| # | 位置 | 一旦发生 | +|---|---|---| +| 1 | `portfolio.positions_view` 取不到现价时把安全垫按 0 记 | 一只实亏 50% 的票被读成「不赚不亏」;更糟的是峰值回吐判据凭空成立,**触发保垫减仓真的卖出去**——而减持方向不设确认门槛 | +| 2 | `proposal_service` 给规则闸的 `day` 是硬编码空壳(`day_chg_from_open: None`) | 「不追高(当日涨幅)」这道闸对**每一笔自主买入从来没真正跑过**。executor 那条路一直取的是真快照,只有这里漏了 | +| 3 | `plan_liquidate_all` 悄悄漏掉取不到现价的票 | 一键清仓是 danger 级命令,语义是「全都卖掉」。少卖一只就是留了个敞口,而 notes 写的是「全部 N 只」——那个 N 数的是下单条数不是持仓只数 | +| 4 | 参数表读不到时 `PMS_GLOBAL_BUY_HALT` 按 `False` 放行 | 数据库一抖,全局暂停买入自己失效。安全开关 fail-open 是方向性错误 | +| 5 | `PMS_RECON_STREAK_YMD` 漏在可写白名单外 | 「按交易日推进」的修复被 `set_param` 静默拒写打回原形,计数退化成按次累加。**单测全绿——因为单测只测纯函数** | +| 6 | `plan_sector_cap` 算错分母 | 卖了 18 万,行业占比仍然 51.2% > 40%,命令报 DONE | +| 7 | 命令撤在途指令只改本端状态,`dispatcher.cancel` 在命令这条路上**一次都没被调用过** | 切 ws 之后,「全局暂停买入」回一句「已撤销 N 条」,而 QMT 侧子单原封不动继续成交 | +| 8 | `signal_service` 在落库**之前**就把当日去重键 `seen.add` 掉 | 落库失败后,这条风控卖出信号当天再也不会被消化——页面只多一行 error,该卖的票就那么留着 | + +修完之后加了一条**防复发的纪律**,因为八条里至少五条是同一个手势造成的: + +> **关键路径禁止丢弃返回值。** `set_param` / `dispatcher.cancel` / `dispatcher.dispatch` / `executor.cancel_instruction` / `bump_once_guards` / `save_neg_streak` / `qmt_repo.update_order` 这类「失败不抛异常、只把失败写在返回值里」的函数,调用点必须接住并处理。 + +这条纪律由 `test_batch10_units.py` 的 **[A] 组**守着——它是全套单测里唯一一条**静态**用例:AST 扫 `app/` 全库,凡是「整条语句就是一次这类调用、返回值没被任何人接住」的写法直接判失败(`await _db(qmt_repo.update_order, ...)` 这种线程池间接调用也认)。名单在 `SOFT_FAIL`,确实可以丢的写进 `ALLOWED` 并**必须带理由**。[A2] 反过来守着扫描器本身:造一个丢返回值的调用,扫不出来就算失败——守卫写错会让 [A1] 永远绿,那比没有守卫更糟。 + +顺带修出来的第九条,是**本轮改动自己引入的**:`positions_view` 改成「取不到现价时拿摊薄成本顶住 `price`」之后,`planner._usable` 光看 `price > 0` 就漏了——拿成本价算出来的市值会被当成真市值去凑「释放 20 万」,凑够了报 DONE 而实际卖出金额对不上。现在 `_usable` 一并排掉 `price_ok is False` 的票。 ### 下一步(按可动工顺序) @@ -316,6 +343,7 @@ make t-gate # 随时: 规则闸/研判闸拒了什么、为什么 | 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~~ | ✅ 2026-07-31:`Makefile`(`make help` 看全部)。`deploy` 包 `scripts/deploy.sh`;另有 `test` / `initdb` / `check` / `probe` / `changes` / `industry` / `ws-status`。一次性命令统一带 `--no-deps`,免得跑个单测把 beat/worker 也拉起来 | +| 13 | ~~静默失败专项:八条全修 + 关键路径禁止丢弃返回值~~ | ✅ 2026-07-31:见上一节。新增 `test_batch10_units.py` 27 例,其中 [A] 组是静态扫描守卫 | ### 运行态注意(2026-07-31 收尾时的状态) @@ -323,6 +351,7 @@ make t-gate # 随时: 规则闸/研判闸拒了什么、为什么 - 账本已清空(`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`。 +- 静默失败专项修完之后,**几处「以前默默过去」的地方现在会明着报失败**,这是设计如此:开关类命令(`HALT_BUY` / `RESUME_ALL` 等)写不进参数就置 `CANCELLED` 并回 `ok=False`,不再显示「已完成」;`daily_settle` 在对账被拦(爆炸半径 / 两源无应答 / 成本价体检不过)或连续不一致升到 ERROR 时 `ok=False`;`premarket` 在「该踩刹车却没踩上」时把它记进 `errors` 而不是 warnings。**看到这些报错先别急着改代码——多半是数据库或下游真的出问题了,以前只是没人告诉你。** - 新增第 15 张表 `pms_plan_snapshot`(名册快照),**部署后要跑一次 `make initdb`**,否则榜单变化那段会一直报「读名册快照失败:表不存在」——候选池不受影响。第一份快照落下之前不报任何变化(`kind=first`),这是设计如此不是坏了。 ### ws 通道实现清单 @@ -395,5 +424,6 @@ docker compose --profile ws up -d pms-ws && docker compose logs -f pms-ws - 配置分两层:基础设施连接串在 `.env`(服务器手工维护,不入库,模板见 `.env.example`);业务参数在 `config/settings.py` 只是初值,上线后经管理页面修改并持久化到 `pms_runtime_param` 表。**业务代码禁止直接读 settings 取业务参数**,一律走 `services/param_store.py`。 > 例外是**密钥**:`PMS_QMT_SIGN_SEED_HEX` / `PMS_QMT_PEER_PUBKEY_B64` 只走 `.env`(协议 §10.1.1),已列入 `param_store.SECRET_KEYS`——页面读不到、改不了、快照里也不出现。新增密钥类配置记得同步加进那个元组,否则它会被当成普通业务参数显示在参数设置页上。 - 新增纯逻辑一律进 `app/core/`(零外部依赖 + 配套单测);需要连库的编排进 `app/services/`,并保证连库失败时降级而非崩页。ws 通道的协议编解码在 `app/core/ws_codec.py`,改动后**必须先跑通协议 §2.1.1 的测试向量**(`test_batch6_units.py` 的 [A] 组)——规范化串两边写不一致的话,联调只会告诉你"签名验不过",看不出差在哪一段。 +- **关键路径禁止丢弃返回值。** 本系统里大量写入函数失败时不抛异常,只回 `{"ok": False, ...}` 或影响 0 行——`param_store.set_param`、`dispatcher.dispatch/cancel`、`executor.cancel_instruction`、`proposal_service.bump_once_guards`、`portfolio.save_neg_streak`、`qmt_repo.update_order` 都是。**写这类调用时返回值必须接住并处理**,否则写入没发生而调用方照常报成功。`test_batch10_units.py` 的 [A] 组 AST 扫全库守着这条;新加一个软失败函数,记得同时加进它的 `SOFT_FAIL` 名单。 - 里程碑(设计定稿、建表、各期上线)及时 git 提交。 - **接手前先读三份**:本文件 → `POSITION_MGMT_DESIGN.md`(V0.4 定稿,**勿改设计**)→ `QMT_WS_PROTOCOL.md`(V1.0 定稿,下发通道的唯一依据)。三份读完即可开工,不需要额外的口头背景。 diff --git a/app/core/action_engine.py b/app/core/action_engine.py index 0afb176..f675c45 100644 --- a/app/core/action_engine.py +++ b/app/core/action_engine.py @@ -197,6 +197,16 @@ def scan(*, positions: list, params: dict, market: dict, skip: set = None) -> di code = p.get("ts_code") if not code or int(p.get("total_qty") or 0) <= 0: continue + # 取不到现价的票整只跳过, **并且留痕**。上游 (portfolio.positions_view) 在拿不到 + # 行情时会用摊薄成本顶住 price 让市值还能算, 但那个价不是行情 —— 拿它评动作会得出 + # 「安全垫恰好 0」「现价恰好等于成本」这类看着正常、实则凭空的结论。 + # 四个 evaluator 目前都会因为 cushion_pct is None 而自然返回 None, 但那是**碰巧** + # 兜住了: eval_fill 还会拿这个假价去比支撑位。所以在入口显式挡掉, 并写进 skipped —— + # 「什么都没发生」和「明着跳过了」在页面上必须是两回事。 + if p.get("price_ok") is False: + skipped.append({"ts_code": code, "action": "*", + "why": "取不到现价 (price 是拿摊薄成本顶的), 本轮不评估该票"}) + continue mkt = (market or {}).get(code) or {} frozen = (p.get("frozen_reason") or "NONE") != "NONE" for action, fn in EVALUATORS: diff --git a/app/core/planner.py b/app/core/planner.py index ea224b8..acf5fe2 100644 --- a/app/core/planner.py +++ b/app/core/planner.py @@ -67,7 +67,13 @@ def _sortable(p, key, reverse=False): def _usable(positions, exclude_codes): - """可纳入减持方案的持仓: 有数量、有现价、不在排除名单 (已有在途方案的票)。""" + """可纳入减持方案的持仓: 有数量、**有真实现价**、不在排除名单 (已有在途方案的票)。 + + `price_ok is False` 的票必须排掉。positions_view 取不到现价时会拿摊薄成本把 price + 顶上 (免得组合市值凭空塌一块), 于是这里光看 `price > 0` 就漏了 —— 拿成本价算出来的 + 市值会被当成真市值去凑"释放 20 万", 凑够了报 DONE, 实际卖出金额对不上。缺价的票要么 + 被排除 (凑额类方案), 要么单独走 need_price 分支 (一键清仓)。 + """ ex = set(exclude_codes or ()) out = [] for p in positions or []: @@ -75,6 +81,8 @@ def _usable(positions, exclude_codes): continue if int(p.get("total_qty") or 0) <= 0 or float(p.get("price") or 0) <= 0: continue + if p.get("price_ok") is False: + continue out.append(p) return out @@ -423,15 +431,36 @@ def plan_liquidate_all(*, positions: list, pending_buys=None) -> dict: b.get("amount") or 0, P_HALT, "一键清仓: 撤销在途买入", tier="1_停新买", cancel_instruction_id=b.get("instruction_id"))) acc = 0.0 - for p in sorted(_usable(positions, ()), key=lambda x: (-mv_of(x), x["ts_code"])): + # **取不到现价的持仓不许静默消失。** `_usable` 会把 price<=0 的票滤掉 (对减持凑额类 + # 方案是对的: 算不出金额就没法凑), 但一键清仓是 danger 级命令, 语义是"全都卖掉"—— + # 少卖一只就是留了个敞口, 而原来的 notes 写的是「全部 N 只持仓清仓」, 那个 N 数的是 + # 下单条数不是持仓数, 读起来像全清了。停牌票在行情库里没有当日分钟线, avg_cost=0 的 + # RECON 批次也会让 price 归 0, 都能走到这里。 + # 照 plan_exit_stock 的口径办: 照样下整票卖单 (金额记 0, 由择时按实时行情定价)。 + priced = _usable(positions, ()) + priced_codes = {p["ts_code"] for p in priced} + unpriced = [p for p in (positions or []) + if int(p.get("total_qty") or 0) > 0 and p["ts_code"] not in priced_codes] + for p in sorted(priced, key=lambda x: (-mv_of(x), x["ts_code"])): qty = int(p.get("total_qty") or 0) amt = qty * float(p["price"]) acc += amt items.append(_item(p["ts_code"], A_EXIT, SIDE_SELL, qty, amt, P_WEAK, "一键清仓 (紧急, 不做择时优化)", tier="全部清仓")) - return {"ok": True, "target_amount": round(acc, 2), "planned_amount": round(acc, 2), - "gap": 0.0, "items": items, "rejects": [], - "notes": [f"全部 {len([i for i in items if i['action'] == A_EXIT])} 只持仓清仓"]} + for p in sorted(unpriced, key=lambda x: x["ts_code"]): + qty = int(p.get("total_qty") or 0) + items.append(_item(p["ts_code"], A_EXIT, SIDE_SELL, qty, 0.0, P_WEAK, + "一键清仓 (紧急): **取不到现价**, 金额待择时按实时行情定", + tier="全部清仓", need_price=True)) + notes_extra = ([f"**{len(unpriced)} 只取不到现价**, 已照常下清仓单但金额算不出: " + + ", ".join(p["ts_code"] for p in unpriced[:6]) + + ("…" if len(unpriced) > 6 else "")] if unpriced else []) + n_exit = len([i for i in items if i["action"] == A_EXIT]) + n_held = len([p for p in (positions or []) if int(p.get("total_qty") or 0) > 0]) + # 「全部 N 只」的 N 必须是**持仓只数**而不是下单条数 —— 两者不等时那句话就是假的 + return {"ok": n_exit >= n_held, "target_amount": round(acc, 2), + "planned_amount": round(acc, 2), "gap": 0.0, "items": items, "rejects": [], + "notes": [f"持仓 {n_held} 只, 已下清仓单 {n_exit} 只"] + notes_extra} def plan_sector_exit(*, sector: str, positions: list) -> dict: @@ -460,7 +489,12 @@ def plan_sector_cap(*, sector: str, cap: float, positions: list) -> dict: if port_mv <= 0 or sec_mv <= 0: return {"ok": True, "items": [], "target_amount": 0.0, "planned_amount": 0.0, "gap": 0.0, "notes": [f"行业[{sector}]无持仓, 仅写入上限参数"], "rejects": []} - over = sec_mv - float(cap) * port_mv + # 卖出会**同时**缩小分子和分母, 所以要解 (sec-x)/(port-x) = cap → x = (sec-cap*port)/(1-cap)。 + # 原来写的是 sec - cap*port, 少卖 1/(1-cap) 倍 (cap=40% 时少 1.67 倍): 实测组合 100 万、 + # 银行 60 万、上限 40% → 方案卖 18 万, 卖完还是 51.2%, 而命令报 ok=True 判 DONE。 + # 三处口径就此打架: 命令说完成、页面行业条还是红的、sizer 的 SECTOR_RATIO 继续拦买入。 + cap_f = min(0.999999, max(0.0, float(cap))) + over = (sec_mv - cap_f * port_mv) / (1.0 - cap_f) if over <= 0: return {"ok": True, "items": [], "target_amount": 0.0, "planned_amount": 0.0, "gap": 0.0, "notes": [f"行业[{sector}]占比 {sec_mv / port_mv:.1%} 未超上限 " @@ -475,9 +509,19 @@ def plan_sector_cap(*, sector: str, cap: float, positions: list) -> dict: acc += amt items.append(_item(p["ts_code"], A_TRIM, SIDE_SELL, q, amt, P_PRORATA, f"行业[{sector}]超限, 等比减持 {q} 股", tier="行业限额")) - return {"ok": acc > 0, "target_amount": round(over, 2), "planned_amount": round(acc, 2), - "gap": round(max(0.0, over - acc), 2), "items": items, "rejects": [], - "notes": [f"行业[{sector}]占比 {sec_mv / port_mv:.1%} → 需减 {over:,.0f} 元"]} + gap = max(0.0, over - acc) + after = (sec_mv - acc) / (port_mv - acc) if port_mv - acc > 0 else 0.0 + notes = [f"行业[{sector}]占比 {sec_mv / port_mv:.1%} → 需减 {over:,.0f} 元, " + f"方案减 {acc:,.0f} 元 → 减后 {after:.1%}"] + # **减完还超限就要明说。** 一手取整让每只票只能整手卖, 残留几乎必然存在; 原来只报 + # 「需减 X 元」而不报减后是多少, 命令照样判完成, 页面行业条还红着却没人知道为什么。 + if after > cap_f + 1e-9: + notes.append(f"**减后仍超上限 {cap_f:.0%}** (一手取整所致, 差 {gap:,.0f} 元) —— " + f"该行业的买入会继续被规则闸拦住, 要彻底降下来需再下一次命令或手工减") + return {"ok": acc > 0 and after <= cap_f + 1e-9, + "target_amount": round(over, 2), "planned_amount": round(acc, 2), + "gap": round(gap, 2), "sector_ratio_after": round(after, 4), + "items": items, "rejects": [], "notes": notes} def plan_halt_buy(*, pending_buys: list) -> dict: @@ -488,7 +532,10 @@ def plan_halt_buy(*, pending_buys: list) -> dict: for b in (pending_buys or [])] return {"ok": True, "items": items, "target_amount": 0.0, "planned_amount": 0.0, "gap": 0.0, "rejects": [], - "notes": [f"撤销在途买入指令 {len(items)} 条"] if items else ["无在途买入指令"]} + # 措辞是"点名"不是"已撤销": planner 只负责挑出要撤哪些, 真撤掉几条要等 + # command_service 调 executor.cancel_instruction 回来才知道。写成"已撤销 N 条" + # 就是在承诺一件还没发生的事 —— 下游拒了照样显示 N 条。 + "notes": [f"点名撤销在途买入指令 {len(items)} 条"] if items else ["无在途买入指令"]} def plan_halt_all(*, pending_instructions: list) -> dict: @@ -500,7 +547,7 @@ def plan_halt_all(*, pending_instructions: list) -> dict: for b in (pending_instructions or [])] return {"ok": True, "items": items, "target_amount": 0.0, "planned_amount": 0.0, "gap": 0.0, "rejects": [], - "notes": [f"撤销在途指令 {len(items)} 条"] if items else ["无在途指令"]} + "notes": [f"点名撤销在途指令 {len(items)} 条"] if items else ["无在途指令"]} # ================================================================ diff --git a/app/core/rule_gate.py b/app/core/rule_gate.py index 9f42782..e7bbf12 100644 --- a/app/core/rule_gate.py +++ b/app/core/rule_gate.py @@ -114,7 +114,13 @@ def _check_buy(failed, warns, qty, price, ctx, pos, day, prm, flg, is_cmd, hard) # 不追高: 当日涨幅 与 距 MA5 两道 dayup = day.get("day_chg_from_open") cap_dayup = _num(prm.get("buy_halt_dayup"), 0.05) - if dayup is not None and _num(dayup) > cap_dayup: + if dayup is None: + # 缺输入**必须留痕**, 与下面 MA5_MISSING 同口径。原来只有 `is not None` 一个条件, + # 缺了就整道闸无声跳过 —— 而自主提议那条路一直硬编码 None, 于是动作引擎产出的每一笔 + # 买单都没过过这道闸, 返回值还是 `passed=True, failed=[], warnings=[]`, 与 + # 「涨幅 2%、检查通过」一字不差, 账本里也查不出来。 + warns.append("DAYUP_MISSING: 取不到当日涨幅, 不追高(涨幅)一项未校验") + elif _num(dayup) > cap_dayup: failed.append(f"NO_CHASE_DAYUP: 当日涨幅 {_num(dayup):.2%} > 上限 {cap_dayup:.0%}") ma5 = _num(day.get("ma5")) cap_ma5 = _num(prm.get("no_chase_ma5"), 0.06) diff --git a/app/services/command_service.py b/app/services/command_service.py index 7f6642c..30a7db9 100644 --- a/app/services/command_service.py +++ b/app/services/command_service.py @@ -232,23 +232,44 @@ def plan_command(cmd: dict) -> dict: return {"ok": False, "command_id": command_id, "errors": [f"DB_ERROR: {e}"]} # 撤单类动作立刻执行 (撤销在途买入/全部在途指令) - cancelled = _cancel_marked_instructions(items) + cancel_r = _cancel_marked_instructions(items) + cancelled = cancel_r["cancelled"] + cancel_failed = cancel_r["failed"] - # 开关类命令写运行参数 + # 开关类命令写运行参数。 + # **这行返回值绝对不能丢**: HALT_BUY / HALT_ALL 的"刹车"就是这一次写入 —— 规则闸 + # 每一跳都去 param_store 读 PMS_GLOBAL_BUY_HALT / PMS_GLOBAL_EXEC_HALT, 参数没写进去 + # 就等于刹车没踩。而这类命令 spec.instant=True, 原来无论写没写成都直接置 DONE 并回 + # ok=True —— 页面上"全局暂停买入 已完成", 买单照下。宁可让命令报失败, 也不能让用户 + # 以为已经停了。(2026-07-31 静默失败专项) + switch_err = None if spec.get("switch_key") is not None: - param_store.set_param(spec["switch_key"], spec.get("switch_value"), "command") + w = param_store.set_param(spec["switch_key"], spec.get("switch_value"), "command") or {} + if not w.get("ok"): + switch_err = f"{spec['switch_key']} 写入失败: {w.get('error')}" + logger.error("[命令] %s %s 开关参数没写进去 —— **该命令没有生效**: %s", + command_id, cmd_type, switch_err) rejects = result.get("rejects") or [] + notes = list(result.get("notes", [])) + if switch_err: + notes.insert(0, f"⚠ 开关未生效: {switch_err}") + # 撤单没撤干净必须顶到 notes 最前面 —— planner 写 notes 的时候还不知道撤单结果, + # 它数的是"点名了几条", 真撤掉几条只有这里知道。埋在 cancel_failed 里没人看。 + if cancel_failed: + notes.insert(0, f"⚠ {len(cancel_failed)} 条在途指令**没撤掉, 仍在下游挂着**: " + + "; ".join(f"{x['instruction_id']}({x['why']})" for x in cancel_failed[:5])) progress = {"target_amount": result.get("target_amount", 0.0), "planned_amount": result.get("planned_amount", 0.0), "done_amount": 0.0, "gap": result.get("gap", 0.0), "plan_count": len(rows), "deadline": str(deadline), - "notes": result.get("notes", []), "rejects": rejects, + "notes": notes, "rejects": rejects, # 一行话说清"为什么只有这么少 / 一条都没有"。planner 的 notes 只会说 # "候选与补仓空间不足, 缺口 X 元" —— 那读起来像"没票可买", 而真相往往是 # 有一堆候选、全被同一道闸拒了。不聚合出来就得去翻 rejects 原始清单。 "reject_summary": pl.summarize_rejects(rejects), - "cancelled_instructions": cancelled} + "cancelled_instructions": cancelled, + "cancel_failed": cancel_failed} # 零方案的两种情形必须分开 (2026-07-31 修) # ------------------------------------------------------------------ @@ -258,14 +279,23 @@ def plan_command(cmd: dict) -> dict: # 拒绝撤销, 用户既看不出没执行成、也退不回来。日报统计也会把它算成完成的命令。 # 次序要紧: `instant` 是**命令规格**的属性 (下达即完成), 优先于有没有方案 —— # HALT_BUY 这类开关命令照样会产出撤单动作 (rows 非空), 但它下达完就该是 DONE。 - if spec.get("instant"): + # instant 命令自己失败了也不许算 DONE。`_adjust_window` / `_cancel_target` 靠返回 + # result["ok"]=False 报"目标命令不存在 / 该命令不可撤销", 原来这个 ok 全库没人读 —— + # 用户看到"撤销命令 已完成", 而目标命令还在跑。开关写不进去同理。 + note, errors = None, [] + instant_failed = switch_err or (spec.get("instant") and result.get("ok") is False) + if instant_failed: + status = cs.ST_CANCELLED + note = ("命令未生效: " + (switch_err or "; ".join( + str(x) for x in (result.get("notes") or ["方案生成器报告失败"]))))[:280] + errors = [note] + logger.error("[命令] %s %s 未生效 → 置 CANCELLED。%s", command_id, cmd_type, note) + elif spec.get("instant"): status = cs.ST_DONE elif rows: status = cs.ST_EXECUTING else: status = cs.ST_CANCELLED - note = None - if status == cs.ST_CANCELLED: note = ("未产出任何方案: " + (progress["reject_summary"] or "候选池为空"))[:280] logger.warning("[命令] %s %s 未产出任何方案 → 置 CANCELLED。%s", command_id, cmd_type, note) @@ -274,8 +304,8 @@ def plan_command(cmd: dict) -> dict: if status in (cs.ST_DONE, cs.ST_CANCELLED) else None)) _ledger_rejects(rejects, command_id) - return {"ok": True, "command_id": command_id, "status": status, "plan": progress, - "items": items, "errors": []} + return {"ok": not errors, "command_id": command_id, "status": status, "plan": progress, + "items": items, "errors": errors} def _dispatch_planner(cmd_type: str, p: dict, cmd: dict) -> dict: @@ -310,8 +340,17 @@ def _dispatch_planner(cmd_type: str, p: dict, cmd: dict) -> dict: if cmd_type == "SECTOR_EXIT": return pl.plan_sector_exit(sector=p["sector"], positions=positions) if cmd_type == "SECTOR_CAP": - param_store.set_param(f"PMS_SECTOR_CAP_{p['sector']}", float(p["cap"]), "command") - return pl.plan_sector_cap(sector=p["sector"], cap=float(p["cap"]), positions=positions) + r = pl.plan_sector_cap(sector=p["sector"], cap=float(p["cap"]), positions=positions) + # 这条命令有两半: 减到线内 (方案) + 把线记下来 (参数)。参数没写上就只减了这一次, + # 上限并没有立住 —— 必须说出来, 不能默默只做一半。 + w = param_store.set_param(f"PMS_SECTOR_CAP_{p['sector']}", + float(p["cap"]), "command") or {} + if not w.get("ok"): + logger.error("[命令] 行业上限参数没写进去 %s: %s", p["sector"], w.get("error")) + r = {**r, "notes": [f"⚠ 本次已减到线内, 但行业上限**没能记下来** " + f"({w.get('error')}), 下一轮不会自动守这条线"] + + list(r.get("notes") or [])} + return r if cmd_type == "OPEN_TARGET": code = p["ts_code"] @@ -441,18 +480,38 @@ def _pending_instructions() -> list: "side": r.get("side")} for r in rows] -def _cancel_marked_instructions(items: list) -> list: - out = [] +def _cancel_marked_instructions(items: list) -> dict: + """撤销命令点名的在途指令。返回 {"cancelled": [...], "failed": [{id, why}]}。 + + **必须走 executor.cancel_instruction, 不能只把本端的行标成 CANCELLED。** + 2026-07-31 查出来: 原来这里只 `update_instruction(status="CANCELLED")`, 全仓库 + `dispatcher.cancel` 在命令这条路上**一次都没被调用过**。切到 ws 之后, 「全局暂停买入」 + 会回一句「已撤销 N 条」, 而 `pms_qmt_order` 里的子单原封不动继续挂着、继续成交 —— + 用户以为踩了刹车, 实际只是本端账面上把它划掉了。HALT_BUY / HALT_ALL / LIQUIDATE_ALL / + REDUCE_EXPOSURE / 撤销命令 五条路全中。 + + 失败的要**单独列出来**, 不能吞成 warning 后照样报"已撤销 N 条" —— 那个 N 原来数的是 + 命令点名的条数, 不是真撤掉的条数。 + """ + from app.services import executor + ok, failed = [], [] for it in items or []: iid = it.get("cancel_instruction_id") if not iid: continue try: - pms_repo.update_instruction(iid, status="CANCELLED") - out.append(iid) + r = executor.cancel_instruction(iid, reason="命令撤销在途指令") or {} + if r.get("ok"): + ok.append(iid) + else: + failed.append({"instruction_id": iid, + "why": r.get("message") or r.get("error") or "撤销未成功"}) except Exception as e: - logger.warning("撤销指令失败 %s: %s", iid, e) - return out + logger.exception("撤销指令失败 %s", iid) + failed.append({"instruction_id": iid, "why": f"{type(e).__name__}: {e}"}) + if failed: + logger.error("[命令] %s 条在途指令没撤掉, **它们仍在下游挂着**: %s", len(failed), failed) + return {"cancelled": ok, "failed": failed} def _codes_with_live_plans() -> list: diff --git a/app/services/executor.py b/app/services/executor.py index 0daae92..dc9c4eb 100644 --- a/app/services/executor.py +++ b/app/services/executor.py @@ -287,7 +287,19 @@ def cancel_instruction(instruction_id: str, reason: str = "页面人工撤销") return {"ok": False, "error": "指令不存在"} if ins["status"] not in LIVE: return {"ok": False, "error": f"指令处于 {ins['status']}, 不可撤销"} - r = dispatcher.cancel(instruction_id=instruction_id, dispatch_ref=ins.get("dispatch_ref")) + r = dispatcher.cancel(instruction_id=instruction_id, + dispatch_ref=ins.get("dispatch_ref")) or {} + # **下游没撤成就不能在本端标 CANCELLED。** 原来 r 拿到手却从不检查, 一律回「已撤销」 + # 并把父指令置终态 —— 而终态之后没有任何东西再跟踪它, 下游那张委托继续挂着、继续成交, + # 成交回来还会因为找不到在途父指令而变成孤儿。撤不掉就保持在途, 让它继续被择时/收口/ + # 对账看见, 由人再处理。 + if not r.get("ok"): + logger.error("[撤单] 下游未受理 %s: %s —— 指令保持在途, 下游委托可能仍挂着", + instruction_id, r.get("error") or r) + return {"ok": False, "downstream": r, "instruction_id": instruction_id, + "message": f"指令 {instruction_id} **未能撤销**: " + f"{r.get('error') or '下游未受理'} —— 指令仍在途, " + f"请在 QMT 侧确认该委托是否还挂着"} pms_repo.update_instruction(instruction_id, status=ST_CANCELLED) pms_repo.insert_ledger(ts_code=ins["ts_code"], action=ins.get("action"), arbiter="user", verdict="REJECT", price_at=0, ref_id=instruction_id, reason=reason) diff --git a/app/services/ledger_service.py b/app/services/ledger_service.py index b5209f1..dc242b8 100644 --- a/app/services/ledger_service.py +++ b/app/services/ledger_service.py @@ -623,9 +623,18 @@ def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict: st = rc.advance_streak(streak, param_store.get_int(STREAK_YMD_KEY, 0), int(datetime.now().strftime("%Y%m%d")), bool(diffs)) if st["changed"] or st["ymd"] != param_store.get_int(STREAK_YMD_KEY, 0): - # 必须走 ParamStore 写入: 直接写库不会失效缓存, 连续天数会一直读到旧值 - param_store.set_param(STREAK_KEY, st["streak"], "system") - param_store.set_param(STREAK_YMD_KEY, st["ymd"], "system") + # 必须走 ParamStore 写入: 直接写库不会失效缓存, 连续天数会一直读到旧值。 + # **返回值必须接**: set_param 写不进去时只回 {"ok": False, "error": ...} 而不抛 + # 异常 —— 2026-07-31 就是因为把它丢了, STREAK_YMD 键不在白名单里被静默拒写, + # prev_ymd 永远读回 0, 按日推进的修复形同虚设而单测全绿 (单测只测纯函数)。 + for k, v in ((STREAK_KEY, st["streak"]), (STREAK_YMD_KEY, st["ymd"])): + w = param_store.set_param(k, v, "system") + if not (w or {}).get("ok"): + logger.error("[对账] 连续天数写入失败 %s=%s: %s —— " + "按日推进将失效, 计数会退化成按次累加", k, v, + (w or {}).get("error")) + out.setdefault("warnings", []).append( + f"连续天数未能落库 ({k}): {(w or {}).get('error')}") streak = st["streak"] out["streak"] = streak out["severity"] = rc.recon_severity(streak if diffs else 0, @@ -954,6 +963,9 @@ def premarket() -> dict: out["errors"].append(f"{code} 参考位取数失败: {e}") try: out["brake"] = _settle_brake() + # 「该踩没踩上」必须冒到盘前准备的 errors 里 —— 它不是一条 warning, 是风控没生效 + for w in (out["brake"].get("warnings") or []): + (out["errors"] if "刹车未生效" in w else out.setdefault("warnings", [])).append(w) except Exception as e: out["errors"].append(f"刹车结算失败: {e}") out["ok"] = not out["errors"] @@ -961,7 +973,16 @@ def premarket() -> dict: def _settle_brake() -> dict: - """组合刹车: 自高水位回撤 ≥ 阈值 → 自主增持停 N 个交易日 (命令类不受限)。""" + """组合刹车: 自高水位回撤 ≥ 阈值 → 自主增持停 N 个交易日 (命令类不受限)。 + + 两处 set_param 的返回值都必须接 (2026-07-31 静默失败专项): + * PMS_BRAKE_UNTIL 写不进去 = 刹车根本没踩。proposal_service 每一跳重新去 + param_store 读这个键判 brake_active, 读不到就是 False, 自主增持照跑。而本函数 + 原来照样 return {"brake_until": until, "active": True} 并打日志"暂停至 X" —— + 日报、页面、日志三处都说停了, 实际一路买穿整个回撤。 + * PMS_HIGH_WATER 写不进去 = 高水位永远停在旧值, 回撤按陈年高点算, 阈值再也够不着。 + 风控是慢性失效, 没有任何一处会报错。 + """ v = portfolio.positions_view() mv = v["totals"]["portfolio_mv"] hw = param_store.get_float("PMS_HIGH_WATER", 0.0) @@ -969,17 +990,34 @@ def _settle_brake() -> dict: days = param_store.get_int("PMS_BRAKE_DAYS", 3) until = param_store.get_int("PMS_BRAKE_UNTIL", 0) today = td.ymd() + warnings = [] if mv > hw: - param_store.set_param("PMS_HIGH_WATER", mv, "system") - hw = mv + w = param_store.set_param("PMS_HIGH_WATER", mv, "system") or {} + if w.get("ok"): + hw = mv + else: + logger.error("[刹车] 高水位写入失败 (仍按旧值 %s 算回撤): %s", hw, w.get("error")) + warnings.append(f"高水位未能落库, 回撤按旧高点 {hw:,.0f} 计算: {w.get('error')}") drawdown = (1 - mv / hw) if hw > 0 else 0.0 + engaged = False if hw > 0 and drawdown >= dd_limit and today >= until: - until = td.ymd(td.next_trade_day(datetime.now().date(), days)) - param_store.set_param("PMS_BRAKE_UNTIL", until, "system") - logger.warning("[刹车] 自高水位回撤 %.1f%% ≥ %.0f%%, 自主增持暂停至 %s", - drawdown * 100, dd_limit * 100, until) - return {"high_water": hw, "portfolio_mv": mv, "drawdown": round(drawdown, 4), - "brake_until": until, "active": today < until} + nxt = td.ymd(td.next_trade_day(datetime.now().date(), days)) + w = param_store.set_param("PMS_BRAKE_UNTIL", nxt, "system") or {} + if w.get("ok"): + until, engaged = nxt, True + logger.warning("[刹车] 自高水位回撤 %.1f%% ≥ %.0f%%, 自主增持暂停至 %s", + drawdown * 100, dd_limit * 100, until) + else: + logger.error("[刹车] **该踩刹车但没踩上** —— 回撤 %.1f%% ≥ %.0f%%, " + "PMS_BRAKE_UNTIL 写入失败: %s。自主增持仍在放行", + drawdown * 100, dd_limit * 100, w.get("error")) + warnings.append(f"**刹车未生效**: 回撤 {drawdown:.1%} 已达阈值但 " + f"PMS_BRAKE_UNTIL 写不进去 ({w.get('error')}), 自主增持仍在放行") + out = {"high_water": hw, "portfolio_mv": mv, "drawdown": round(drawdown, 4), + "brake_until": until, "active": today < until, "engaged": engaged} + if warnings: + out["warnings"] = warnings + return out def daily_settle() -> dict: @@ -991,7 +1029,18 @@ def daily_settle() -> dict: except Exception as e: out["errors"].append(f"除权检测失败: {e}") try: - out["steps"]["recon"] = reconcile() + r = out["steps"]["recon"] = reconcile() + # reconcile 的"我拒绝对账"是**返回值**不是异常 (爆炸半径拦截 / 两源无应答 / + # 成本价体检不过 → ok=False 或 severity=ERROR)。原来这里只 catch 异常, 于是 + # 一整天没对上账的日子, daily_settle 照样 ok=True 回给调度器和运维页。 + if not r.get("ok"): + out["errors"].append("对账未完成: " + "; ".join( + str(x) for x in (r.get("errors") or [r.get("blocked") or "原因见 recon 步骤"]))) + elif r.get("severity") == "ERROR": + out["errors"].append( + f"对账连续不一致已达 ERROR (连续 {r.get('streak')} 日), 需人工介入") + for w in (r.get("warnings") or []): + out.setdefault("warnings", []).append(f"对账: {w}") except Exception as e: out["errors"].append(f"对账失败: {e}") try: @@ -1027,9 +1076,14 @@ def _settle_cushion() -> dict: cushion_state=x["cushion_state"]) updated += 1 held = {x["ts_code"] for x in v["held"]} - portfolio.save_neg_streak({k: v2 for k, v2 in streak.items() if k in held}) - return {"updated": updated, - "neg_streak": {k: v2 for k, v2 in streak.items() if v2 > 0 and k in held}} + w = portfolio.save_neg_streak({k: v2 for k, v2 in streak.items() if k in held}) or {} + out = {"updated": updated, + "neg_streak": {k: v2 for k, v2 in streak.items() if v2 > 0 and k in held}} + if not w.get("ok"): + # 写不进去要让日终结算整体报失败, 不能只在本函数里留个字段等人去翻 + raise RuntimeError(f"安全垫连负天数未能落库 ({w.get('error')}) —— " + f"降仓命令的「清弱票」判据会停在旧值") + return out # ================================================================ 日报 diff --git a/app/services/param_store.py b/app/services/param_store.py index 0be0dea..19dc609 100644 --- a/app/services/param_store.py +++ b/app/services/param_store.py @@ -43,6 +43,20 @@ RUNTIME_EXTRA = { "PMS_HIGH_WATER": (0.0, float, "组合市值高水位 (刹车判定基准)"), "PMS_REPLAY_CURSOR": ("", str, "成交回放游标 (trading_order.order_id)。留空=首次启动时自动对齐到当前最新成交且不追认历史; 填某个 order_id=从它之后开始; 填 ALL=从头全量回放"), "PMS_RECON_STREAK": (0, int, "连续对账不一致天数"), + # 2026-07-31: 这个键漏在白名单外, 于是 ledger_service 每次写都被 set_param 拒掉 + # (它只回 ok=False 不抛异常, 调用点又把返回值丢了)。后果是 prev_ymd 永远读回 0, + # 「连续 N 日」的按日推进彻底失效, 修复形同虚设。加表时**必须同步加这里**。 + "PMS_RECON_STREAK_YMD": (0, int, "上次推进连续天数的交易日 YYYYMMDD (同日不重复计数)"), +} + +# **读不到时必须按"已暂停"处理的键 (fail-closed)。** +# 这几个开关只活在表里 (settings.py 里根本没有), 所以 `get()` 读不到时会退回 +# RUNTIME_EXTRA 的默认值 False/0 —— 而那恰好等于"放行"。参数表抖一下 (celery worker +# 冷启、连库超时), 全局暂停买入、休假模式、组合刹车就一起静默解除, 且没有任何调用方 +# 看得出这是降级。安全开关的默认方向必须是"拦", 不是"放"。 +FAIL_CLOSED = { + "PMS_GLOBAL_BUY_HALT": True, # 读不到 → 当作已暂停买入 + "PMS_GLOBAL_EXEC_HALT": True, # 读不到 → 当作已暂停执行 } # 页面展示用的中文说明 (settings.py 用行尾注释, pydantic 取不到, 故在此集中维护) @@ -193,6 +207,12 @@ def get(key: str, default=None): return "" rows = refresh() t = _type_of(key) + # 表读不到时, 安全开关按"拦"而不是按默认值放行 —— 见 FAIL_CLOSED 的注释。 + # 只在**表确实读失败**时生效; 表读到了而这个键没设过, 那是真的没设, 照常走默认值。 + if _cache.get("error") and key in FAIL_CLOSED and key not in rows: + logger.error("[参数] 读不到 %s (参数表: %s) —— 按 fail-closed 取 %r, " + "宁可多拦一轮也不放行", key, _cache["error"], FAIL_CLOSED[key]) + return FAIL_CLOSED[key] if key in rows: try: return _coerce(rows[key]["param_value"], t) diff --git a/app/services/portfolio.py b/app/services/portfolio.py index 73c36ee..6596450 100644 --- a/app/services/portfolio.py +++ b/app/services/portfolio.py @@ -75,11 +75,20 @@ def neg_streak_map() -> dict: return {} -def save_neg_streak(m: dict): +def save_neg_streak(m: dict) -> dict: + """写「安全垫连续为负天数」。**必须返回成败, 调用方必须接。** + + 这个计数是 `plan_reduce_exposure` 一档「清弱票」的唯一判据 (连负 N 日 → 优先清)。 + 原来这里把异常吞成 warning 就返回 None, 上层 `_settle_cushion` 照样报 + {"updated": n, ...}、`daily_settle` 照样 ok=True —— 计数永远停在旧值, 降仓命令的 + 第一档静默失灵, 改成卖别的票, 而没有任何一处报错。(2026-07-31 静默失败专项) + """ try: pms_repo.set_param(NEG_STREAK_KEY, json.dumps(m, ensure_ascii=False), "system") + return {"ok": True} except Exception as e: - logger.warning("安全垫连负天数写入失败: %s", e) + logger.error("安全垫连负天数写入失败: %s —— 「清弱票」判据将停在旧值", e) + return {"ok": False, "error": f"{type(e).__name__}: {e}"} def positions_view(*, with_price: bool = True) -> dict: @@ -97,16 +106,27 @@ def positions_view(*, with_price: bool = True) -> dict: code = r["ts_code"] avg_cost = float(r.get("avg_cost") or 0) px = prices.get(code) - if not px or px <= 0: + # 取不到现价时拿摊薄成本顶上, 是为了市值/占比还能算 (否则整张页面塌掉)。 + # **但安全垫绝不能跟着算出来**: px := avg_cost 会让 cushion 恰好等于 0.0, 而 0.0 + # 与"真的不赚不亏"完全同形。后果不是显示难看, 是**凭空产生卖单** —— 一只安全垫 + # 峰值 8% 的停牌票会被判成"回吐至 0%(过半)"而触发 TRIM 保垫减仓, 而保垫减仓 + # **不设确认门槛、直接落指令**; 反过来浮亏 −50% 的票补仓评估档一次都不会触发。 + # 停牌票在行情库里永远没有当日分钟线, 所以这是常规路径不是罕见路径。 + # 价格可以估, 安全垫必须留 None ——「拿不到 ≠ 是 0」。 + price_ok = bool(px and px > 0) + if not price_ok: px = avg_cost or 0.0 if int(r.get("total_qty") or 0) > 0: missing.append(code) qty = int(r.get("total_qty") or 0) mv = qty * px - cp = (px / avg_cost - 1.0) if avg_cost > 0 else None + cp = (px / avg_cost - 1.0) if (price_ok and avg_cost > 0) else None out.append({ "ts_code": code, "status": r.get("status"), "frozen_reason": r.get("frozen_reason"), - "price": round(px, 3), "total_qty": qty, "avail_qty": int(r.get("avail_qty") or 0), + # price_ok=False → 这一行的 price 是拿摊薄成本顶的, 不是行情。动作引擎据此整只 + # 跳过并留痕 (见 action_engine.scan) —— 不是靠 cushion_pct 恰好为 None 兜住。 + "price": round(px, 3), "price_ok": price_ok, + "total_qty": qty, "avail_qty": int(r.get("avail_qty") or 0), "base_qty": int(r.get("base_qty") or 0), "fill_qty": int(r.get("fill_qty") or 0), "add_qty": int(r.get("add_qty") or 0), "dca_qty": int(r.get("dca_qty") or 0), "t0_qty": int(r.get("t0_qty") or 0), "avg_cost": round(avg_cost, 3) or None, diff --git a/app/services/proposal_service.py b/app/services/proposal_service.py index 4780331..c13a200 100644 --- a/app/services/proposal_service.py +++ b/app/services/proposal_service.py @@ -85,9 +85,13 @@ def _route_one(c, view, params, stock_params, brake_active, now, dry_run, out): gate = rule_gate.check( side=side, action=action, qty=c["qty"], price=price, ctx={"ts_code": code, "position": pos, - "day": {"price": price, "vwap": price, - "ma5": (params.get("_mkt") or {}).get(code, {}).get("ma5"), - "day_chg_from_open": None}, + # 当日行情走 _market_ctx 取到的真实快照。原来这里是 + # {"vwap": price, "day_chg_from_open": None} + # 两个都是编的: vwap 拿现价顶等于"现价恰好等于均价", day_chg 给 None 让 + # 「不追高」整道闸静默跳过。取不到就留空, 让规则闸记 *_MISSING 警告。 + "day": {**((params.get("_mkt") or {}).get(code, {}).get("day") or {}), + "price": price, + "ma5": (params.get("_mkt") or {}).get(code, {}).get("ma5")}, "params": {"no_chase_ma5": params.get("no_chase_ma5"), "buy_halt_dayup": params.get("buy_halt_dayup"), "sector_source_ready": view["sector_ready"]}, @@ -163,7 +167,19 @@ def _make_instruction(c, price, now) -> str: progress={"deadline": str(td.window_deadline(now.date(), window)), "is_command": False, "children": [], "auto": True, "reason": c["reason"]}) - bump_once_guards(c["ts_code"], c["action"], c.get("hard_numbers"), now) + g = bump_once_guards(c["ts_code"], c["action"], c.get("hard_numbers"), now) or {} + if not g.get("ok"): + # 指令已经落表了, 不回滚; 但要在评审账本上留一条痕, 否则这条纪律失效没有任何记录 + try: + pms_repo.insert_ledger( + ts_code=c["ts_code"], action=c["action"], arbiter="rule", verdict="WARN", + price_at=float(price or 0), ref_id=iid, + hard_numbers={"once_guard_fields": {k: str(v) for k, v in + (g.get("fields") or {}).items()}}, + reason=f"一次性守卫计数器未写入 ({g.get('error')}) —— " + f"该票 {c['action']} 的「只做一次」本轮失效, 留意重复出手") + except Exception: + logger.exception("一次性守卫失败留痕也没写上 %s", c["ts_code"]) return iid @@ -190,12 +206,21 @@ def bump_once_guards(ts_code: str, action: str, hard_numbers=None, now=None): elif action == "DCA": fields["dca_count"] = int((hard_numbers or {}).get("stage") or 1) if not fields: - return + return {"ok": True, "fields": {}} try: - pms_repo.update_position(ts_code, **fields) + n = pms_repo.update_position(ts_code, **fields) except Exception as e: - # 写不上只是让这条纪律退回原样 (可能重复提), 不该把已落表的指令带崩 - logger.warning("一次性守卫计数器写入失败 %s %s: %s", ts_code, action, e) + # 写不上不该把已落表的指令带崩, 但**必须回成败**: 计数器没写上, 这条一次性纪律 + # 就退回失效状态, 同一票同一动作在首条指令离开 LIVE 之后会再来一次。 + # (2026-07-31 静默失败专项: 原来这里 return None, 调用方拿不到任何信号) + logger.error("一次性守卫计数器写入失败 %s %s: %s —— 该票这条一次性纪律本轮失效", + ts_code, action, e) + return {"ok": False, "fields": fields, "error": f"{type(e).__name__}: {e}"} + if n == 0: + logger.error("一次性守卫计数器没写到任何行 %s %s (持仓行不存在?) —— " + "该票这条一次性纪律本轮失效", ts_code, action) + return {"ok": False, "fields": fields, "error": "持仓行不存在, 影响 0 行"} + return {"ok": True, "fields": fields} def _make_proposal(c, price, verdict) -> str: @@ -236,6 +261,16 @@ def _market_ctx(held: list, now) -> dict: for p in held or []: code = p["ts_code"] d = {"ma5": market.get_ma5(code), "high5": market.get_high5(code)} + # 当日行情快照 —— **规则闸的「不追高(当日涨幅)」全靠它**。 + # 原来这里没取, `_route_one` 给规则闸硬编码 `day_chg_from_open: None`, 而规则闸对 + # None 是整道跳过 → 动作引擎产出的每一笔自主买单 (FILL/ADD/DCA) 从来没过过这道闸, + # 账本里也查不到痕迹。executor 那条路一直是拿 day_snapshot 的, 只有这里漏了。 + # 取不到就留空, 由规则闸记 DAYUP_MISSING 警告 —— 不再无声无息。 + try: + d["day"] = market.day_snapshot(code) or {} + except Exception as e: # 行情不可用不该让整轮扫描崩掉 + logger.warning("[提议] 取当日快照失败 %s: %s", code, e) + d["day"] = {} opened, last_add = p.get("opened_date"), p.get("last_add_date") d["tdays_since_open"] = _tdays_between(opened, today) d["tdays_since_last_add"] = _tdays_between(last_add, today) diff --git a/app/services/signal_service.py b/app/services/signal_service.py index 2def840..89277d8 100644 --- a/app/services/signal_service.py +++ b/app/services/signal_service.py @@ -161,12 +161,18 @@ def _handle(sig, view, prm, seen, ymd, dry_run, out): {**brief, "dry_run": True}) return - seen.add(key) + # 去重键必须**写成功之后**才落。原来是先 seen.add(key) 再落库, 一旦 _make_exit / + # _make_proposal 抛异常 (DB 抖一下、下游拒一次), 上层 catch 住记进 out["errors"], + # 但这一天的去重键已经烧掉了 —— 同一条风控卖出信号后面再来多少次都被当成重复丢弃, + # 指令一条都不会落。失败长得像成功: 页面只多一行 error, 而该卖的票就那么留着了。 + # 2026-07-31 修。 if act == sr.ACT_EXIT: iid = _make_exit(code, d, pos) + seen.add(key) out["exits"].append({**brief, "instruction_id": iid}) else: pid = _make_proposal(code, d, pos, sig) + seen.add(key) out["proposals"].append({**brief, "proposal_id": pid}) diff --git a/app/web/main.py b/app/web/main.py index f7ae42e..0cc7966 100644 --- a/app/web/main.py +++ b/app/web/main.py @@ -261,7 +261,12 @@ def api_decide(proposal_id: str, payload: dict = Body(default={})): # 与自主执行同一口径: 采纳即算「做过一次」, 计数器要跟着走 # (否则页面采纳的那条动作绕开了 §6 的一次性约束) from app.services import proposal_service - proposal_service.bump_once_guards(p["ts_code"], p["action"], hn) + g = proposal_service.bump_once_guards(p["ts_code"], p["action"], hn) or {} + if not g.get("ok"): + # 指令已落表, 不回滚; 但要让页面看见"这条一次性纪律本轮没锁上" + return {"ok": True, "decision": decision, "instruction_id": instruction_id, + "warning": f"一次性守卫计数器未写入 ({g.get('error')}) —— " + f"{p['ts_code']} 的 {p['action']}「只做一次」本轮失效"} return {"ok": True, "decision": decision, "instruction_id": instruction_id} return ok(_decide) diff --git a/app/ws/runner.py b/app/ws/runner.py index e3ca092..a4d815d 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -539,14 +539,30 @@ class WsRunner: logger.error(msg) self._stat["ack_seq_degraded"] = True + def _note_orphan(self, n, type_: str, iid: str): + """update_order 影响 0 行 = 这条委托在出口表里根本没有对应行。 + + 它**不抛异常**, 返回值一丢就什么都没发生过: 本地那张单永远停在在途状态, 一直被 + 算进 queue_depth 与 next_cancel_requests, channel_status 报着并不存在的在途单。 + 计数挂进 _stat, 随心跳落到 pms_ws_state.stat, 页面运维抽屉直接看得到。 + (2026-07-31 静默失败专项) + """ + if n: + return + self._stat["orphan_updates"] = self._stat.get("orphan_updates", 0) + 1 + logger.error("[%s] %s 在 pms_qmt_order 里没有对应行, 状态没落上 " + "(累计孤儿更新 %s 条) —— 本地委托会一直停在在途状态, 请核对出口表", + type_, iid, self._stat["orphan_updates"]) + async def _apply_side_effects(self, type_: str, pl: dict, env: dict): """把上行消息落到 pms_qmt_order 的状态上。**只动通道状态, 不动账本**。""" iid = pl.get("instruction_id") or env.get("corr_id") try: if type_ == wsc.T_ACK: - await _db(qmt_repo.update_order, iid, - status=str(pl.get("status") or wsc.ST_ACCEPTED).upper(), - broker_order_id=pl.get("broker_order_id")) + n = await _db(qmt_repo.update_order, iid, + status=str(pl.get("status") or wsc.ST_ACCEPTED).upper(), + broker_order_id=pl.get("broker_order_id")) + self._note_orphan(n, type_, iid) if pl.get("duplicate"): logger.info("[ack] %s 幂等命中 (对端已受理过), 当前状态 %s", iid, pl.get("status")) @@ -572,10 +588,11 @@ class WsRunner: n, json.dumps(pl, ensure_ascii=False)[:500]) self._maybe_degrade_ack(pl) else: - await _db(qmt_repo.update_order, iid, status=qmt_repo.OS_REJECTED, - reject_code=code[:32], - reject_reason=str(pl.get("reason") or "")[:300], - final_at=datetime.now()) + n = await _db(qmt_repo.update_order, iid, status=qmt_repo.OS_REJECTED, + reject_code=code[:32], + reject_reason=str(pl.get("reason") or "")[:300], + final_at=datetime.now()) + self._note_orphan(n, type_, iid) # 设计 §13: 指令下发失败不自动重发。可重试码也只是记下来, # 由 executor 下一跳按最新行情重新决定 —— 换价重发的判断权在择时模块。 lvl = logger.error if code == "INTERNAL" else logger.warning @@ -645,7 +662,8 @@ class WsRunner: "撤单可能在到期后才到, 按到期处理", iid) logger.info("[order_update] %s 终态 %s, 累计成交 %s 股", iid, status, fields["cum_qty"]) - await _db(qmt_repo.update_order, iid, **fields) + self._note_orphan(await _db(qmt_repo.update_order, iid, **fields), + wsc.T_ORDER_UPDATE, iid) # ============================================================ 确认 async def _ack_loop(self, seed: str): diff --git a/scripts/reset_ledger.py b/scripts/reset_ledger.py index 4e476c2..da1ed31 100644 --- a/scripts/reset_ledger.py +++ b/scripts/reset_ledger.py @@ -175,13 +175,22 @@ def main(): print(f" FAIL 游标重置: {type(e).__name__}: {e}") failed.append(CURSOR_KEY) - # 对账连续不一致天数归零 —— 账本刚清空, 上一轮攒下的 streak 会让下一次对账直接判 ERROR + # 对账连续不一致天数归零 —— 账本刚清空, 上一轮攒下的 streak 会让下一次对账直接判 ERROR。 + # set_param 写不进去是**返回 ok=False 而不抛异常**, 只 try/except 的话归零失败照样 + # 打印 "OK ... 归零", 而下一次对账带着旧 streak 直接升 ERROR。 try: from app.services import param_store - param_store.set_param("PMS_RECON_STREAK", 0, "reset_ledger") - print(" OK PMS_RECON_STREAK 归零 (否则下次对账会带着旧的连续天数直接升 ERROR)") + for _k in ("PMS_RECON_STREAK", "PMS_RECON_STREAK_YMD"): + w = param_store.set_param(_k, 0, "reset_ledger") or {} + if w.get("ok"): + print(f" OK {_k} 归零") + else: + print(f" FAIL {_k} 归零失败: {w.get('error')} —— " + f"下次对账会带着旧的连续天数直接升 ERROR") + failed.append(_k) except Exception as e: - print(f" WARN 对账连续天数归零失败: {type(e).__name__}: {e}") + print(f" FAIL 对账连续天数归零失败: {type(e).__name__}: {e}") + failed.append("PMS_RECON_STREAK") if args.reset_ws: try: diff --git a/scripts/run_tests.py b/scripts/run_tests.py index 807dcec..c129fea 100644 --- a/scripts/run_tests.py +++ b/scripts/run_tests.py @@ -14,8 +14,9 @@ test_batch7_units.py 上游选股计划: 解析/新鲜度/候选筛选/取数守卫 (32 例) test_batch8_units.py 榜单变化: 名册指纹/三种语义/尾部闸/落库往返 (60 例) test_batch9_units.py 成本价体检 / 对账按日推进 / 行业闸 / 取整记账 (47 例) + test_batch10_units.py 静默失败专项: 关键路径不许丢返回值 + 八条实例 (27 例) test_wiring.py 装配自检: 服务层→核心→落表 全链路 (内存桩) (58 例) - 共 351 例 + 共 378 例 任一子集失败即整体失败 (退出码 1)。 """ import os @@ -27,7 +28,7 @@ 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_batch8_units.py", "test_batch9_units.py", - "test_wiring.py"] + "test_batch10_units.py", "test_wiring.py"] def main(): diff --git a/scripts/test_batch10_units.py b/scripts/test_batch10_units.py new file mode 100644 index 0000000..a2bf4f7 --- /dev/null +++ b/scripts/test_batch10_units.py @@ -0,0 +1,623 @@ +# -*- coding: utf-8 -*- +""" +第十批单测: 静默失败专项 (零外部依赖, 不连库) +================================================= +运行: python scripts/test_batch10_units.py + +2026-07-31 专项排查的产物。这一批守的不是某一条业务规则, 而是一类**故障形态**: + + 失败长得像成功。 + +本系统里大量写入函数「失败不抛异常, 只回 {"ok": False, ...} 或 0 行」。返回值一丢, +写入没发生, 而调用方照常往下走、照常回 ok=True、页面照常显示"已完成"。八条实例: + + 1 取不到现价时安全垫按 0 记 → 凭空触发保垫减仓 (卖出方向不设确认门槛, 直接出手) + 2 不追高闸的当日涨幅恒为 None → 自主买入这道闸从来没真正跑过 + 3 清仓命令悄悄漏掉无价的票, 嘴上还说"全部 N 只" + 4 参数表读不到时 HALT 开关按 False 放行 (fail-open) + 5 连续天数的"上次推进日"键漏在白名单外 → 按日推进形同虚设 + 6 行业减仓算错分母 → 卖完仍超限却报 DONE + 7 命令撤在途指令只改本端状态, 下游子单原封不动继续成交 + 8 信号去重键在落库**之前**就烧掉 → 落库失败后这条信号当天再也不会重来 + +[A] 组是这批里唯一一条**静态**用例: 它不测行为, 它扫源码, 守住"关键路径不许丢返回值" +这条纪律本身 —— 新加的调用点一旦又把返回值丢了, 这里立刻红。 +""" +import ast +import os +import sys + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) + +CASES = [] + + +def case(name): + def deco(fn): + CASES.append((name, fn)) + return fn + return deco + + +ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) + + +# ================================================================ +# [A] 关键路径禁止丢弃返回值 (静态扫描) +# ================================================================ +# 这些函数**写失败时不抛异常**, 只把失败写在返回值里。调用点必须接住。 +# key = 函数名, value = 允许的调用者模块名 (None = 不限, 用于 _db(f, ...) 这种间接调用) +SOFT_FAIL = { + "set_param": {"param_store"}, # → {"ok": False, "error": ...} + "save_neg_streak": {"portfolio"}, # → {"ok": False, "error": ...} + "bump_once_guards": {"proposal_service"}, # → {"ok": False, "error": ...} + "cancel_instruction": {"executor"}, # → {"ok": False, "message"|"error": ...} + "dispatch": {"dispatcher"}, # → {"ok": False, "error": ...} + "cancel": {"dispatcher"}, # → {"ok": False, "error": ...} + "update_order": {"qmt_repo"}, # → 影响行数, 0 = 出口表里没这行 +} + +# 确实可以丢的调用点写在这里, **必须带理由**。空着比乱加强。 +ALLOWED = { + # "app/xxx.py:123": "理由", +} + + +def _discarded_calls(path: str) -> list: + """找出「整条语句就是一次调用、返回值没被任何人接住」的软失败调用。""" + with open(path, encoding="utf-8") as f: + tree = ast.parse(f.read(), path) + out = [] + for node in ast.walk(tree): + if not isinstance(node, ast.Expr): # 表达式语句 = 返回值直接扔掉 + continue + call = node.value + if isinstance(call, ast.Await): + call = call.value + if not isinstance(call, ast.Call): + continue + fn = call.func + cand = [] + if isinstance(fn, ast.Attribute): + cand.append((getattr(fn.value, "id", None), fn.attr)) + elif isinstance(fn, ast.Name) and fn.id == "_db" and call.args: + # await _db(qmt_repo.update_order, ...) —— 线程池里跑的同一件事 + a = call.args[0] + if isinstance(a, ast.Attribute): + cand.append((getattr(a.value, "id", None), a.attr)) + for mod, attr in cand: + if attr in SOFT_FAIL and (mod is None or mod in SOFT_FAIL[attr]): + out.append((node.lineno, f"{mod}.{attr}")) + return out + + +@case("[A1] 关键路径禁止丢弃返回值: app/ 全库无「调用了却不看结果」的软失败写入") +def _(): + bad = [] + for d, _dirs, files in os.walk(os.path.join(ROOT, "app")): + for f in sorted(files): + if not f.endswith(".py"): + continue + p = os.path.join(d, f) + rel = os.path.relpath(p, ROOT) + for lineno, what in _discarded_calls(p): + if ALLOWED.get(f"{rel}:{lineno}"): + continue + bad.append(f"{rel}:{lineno} 丢弃了 {what}() 的返回值") + assert not bad, ("以下调用点把「写失败」的返回值扔了 —— 写不进去时调用方一无所知, " + "会照常报成功:\n " + "\n ".join(bad)) + + +@case("[A2] 扫描器本身有效: 造一个丢返回值的调用, 必须扫得出来") +def _(): + # 守住守卫。扫描器写错 (比如 AST 节点类型判断反了) 会让 A1 永远绿, 那比没有更糟。 + import tempfile + src = ("def f(param_store, dispatcher, x):\n" + " param_store.set_param('K', 1, 'test')\n" # ← 该被抓 + " r = param_store.set_param('K', 2, 'test')\n" # ← 接住了, 放行 + " if dispatcher.dispatch(x)['ok']:\n" # ← 用上了, 放行 + " return r\n") + with tempfile.NamedTemporaryFile("w", suffix=".py", delete=False, + encoding="utf-8") as fh: + fh.write(src) + tmp = fh.name + try: + hits = _discarded_calls(tmp) + assert [h[1] for h in hits] == ["param_store.set_param"], hits + assert hits[0][0] == 2, hits + finally: + os.unlink(tmp) + + +@case("[A3] 软失败清单没漏: 名单里的函数确实是「不抛异常只回 ok=False」") +def _(): + from app.services import param_store, portfolio, proposal_service + # set_param: 不可修改的键 → 回 ok=False 而不是抛 + r = param_store.set_param("PROXY_DB_URL", "x") + assert isinstance(r, dict) and r.get("ok") is False and r.get("error"), r + # 另外两个的失败分支在 [E]/[F] 组用桩验证, 这里只钉住"返回的是 dict 不是 None" + assert portfolio.save_neg_streak.__doc__, "save_neg_streak 必须写明返回成败" + assert proposal_service.bump_once_guards.__doc__ + + +# ================================================================ +# [B] 开关命令: 参数没写进去就不许报"已完成" +# ================================================================ +@case("[B1] HALT_BUY 参数写入失败 → 命令置 CANCELLED 并回 ok=False (不能报已完成)") +def _(): + from test_wiring import install_fakes + from app.services import command_service, param_store + from app.core import command_spec as cs + install_fakes(prices={}) + orig = param_store.set_param + try: + param_store.set_param = lambda k, v, by="user": {"ok": False, "error": "库挂了"} + r = command_service.issue("HALT_BUY", {}, issued_by="test") # issue 内部即时排方案 + # 刹车没踩上, 就不许说"已完成" + assert r["ok"] is False, r + assert r["status"] == cs.ST_CANCELLED, r + assert "库挂了" in (r["errors"][0] if r["errors"] else ""), r + assert any("开关未生效" in str(n) for n in r["plan"]["notes"]), r["plan"]["notes"] + finally: + param_store.set_param = orig + + +@case("[B2] 参数写得进去时 HALT_BUY 照常 DONE (修复不能把正常路堵死)") +def _(): + from test_wiring import install_fakes + from app.services import command_service, param_store + from app.core import command_spec as cs + fake = install_fakes(prices={}) + r = command_service.issue("HALT_BUY", {}, issued_by="test") + assert r["ok"] and r["status"] == cs.ST_DONE, r + assert param_store.get("PMS_GLOBAL_BUY_HALT") is True, fake.params + + +@case("[B3] instant 命令自身失败 (撤销一个不存在的命令) → CANCELLED 而不是 DONE") +def _(): + from test_wiring import install_fakes + from app.services import command_service + from app.core import command_spec as cs + install_fakes(prices={}) + r = command_service.issue("CANCEL_COMMAND", {"target_command_id": "CMD_NOPE_0001"}, + issued_by="test") + assert r["ok"] is False and r["status"] == cs.ST_CANCELLED, r + assert "不存在" in str(r["errors"]), r + + +# ================================================================ +# [C] 撤在途指令: 下游拒了必须露出来 +# ================================================================ +@case("[C1] 撤单被下游拒 → 不计入 cancelled, 顶到 notes 第一条") +def _(): + from test_wiring import install_fakes + from app.services import command_service, executor + install_fakes(prices={}) + calls = [] + + def fake_cancel(iid, reason=None): + calls.append(iid) + return {"ok": False, "message": f"下游拒绝撤单: {iid}"} + + orig = executor.cancel_instruction + try: + executor.cancel_instruction = fake_cancel + items = [{"ts_code": "600000.SH", "action": "HALT", "cancel_instruction_id": "INS_A"}, + {"ts_code": "000001.SZ", "action": "HALT", "cancel_instruction_id": "INS_B"}] + r = command_service._cancel_marked_instructions(items) + assert calls == ["INS_A", "INS_B"], calls # 真的调了下游, 不是只改本端 + assert r["cancelled"] == [], r + assert [x["instruction_id"] for x in r["failed"]] == ["INS_A", "INS_B"], r + assert "下游拒绝撤单" in r["failed"][0]["why"], r + finally: + executor.cancel_instruction = orig + + +@case("[C2] 部分成功: cancelled 数的是真撤掉的那些, 不是点名的条数") +def _(): + from test_wiring import install_fakes + from app.services import command_service, executor + install_fakes(prices={}) + orig = executor.cancel_instruction + try: + executor.cancel_instruction = ( + lambda iid, reason=None: {"ok": iid == "INS_OK", "message": "no"}) + r = command_service._cancel_marked_instructions( + [{"cancel_instruction_id": "INS_OK"}, {"cancel_instruction_id": "INS_BAD"}]) + assert r["cancelled"] == ["INS_OK"] and len(r["failed"]) == 1, r + finally: + executor.cancel_instruction = orig + + +@case("[C3] 方案生成器的 notes 措辞是「点名撤销」而不是「已撤销」") +def _(): + from app.core import planner as pl + r = pl.plan_halt_buy(pending_buys=[{"ts_code": "600000.SH", "qty": 100, + "instruction_id": "INS_1"}]) + assert r["items"][0]["cancel_instruction_id"] == "INS_1", r + # planner 还不知道撤没撤成, 不许承诺结果 + assert "点名" in r["notes"][0], r["notes"] + + +# ================================================================ +# [D] 组合刹车: 该踩没踩上必须报错 +# ================================================================ +@case("[D1] PMS_BRAKE_UNTIL 写入失败 → 盘前准备报 errors, 不许 ok=True") +def _(): + from test_wiring import install_fakes + from app.services import ledger_service as ls, param_store, portfolio + install_fakes(prices={}) + orig_set, orig_view = param_store.set_param, portfolio.positions_view + try: + # 高水位 100 万, 现值 80 万 → 回撤 20%, 远超默认 5%, 必须踩刹车 + portfolio.positions_view = lambda **kw: { + "held": [], "positions": [], "params": {}, + "totals": {"portfolio_mv": 800000.0}} + param_store.set_param = lambda k, v, by="user": ( + {"ok": True} if k == "PMS_HIGH_WATER" else {"ok": False, "error": "库挂了"}) + orig_get = param_store.get_float + param_store.get_float = lambda k, d=0.0: (1000000.0 if k == "PMS_HIGH_WATER" + else orig_get(k, d)) + try: + b = ls._settle_brake() + finally: + param_store.get_float = orig_get + assert b["engaged"] is False, b # 没踩上就不能说踩上了 + assert any("刹车未生效" in w for w in b.get("warnings", [])), b + finally: + param_store.set_param, portfolio.positions_view = orig_set, orig_view + + +@case("[D2] 高水位写入失败 → 回撤按旧高点算并留 warning, 不静默") +def _(): + from test_wiring import install_fakes + from app.services import ledger_service as ls, param_store, portfolio + install_fakes(prices={}) + orig_set, orig_view = param_store.set_param, portfolio.positions_view + try: + portfolio.positions_view = lambda **kw: { + "held": [], "positions": [], "params": {}, + "totals": {"portfolio_mv": 500000.0}} + param_store.set_param = lambda k, v, by="user": {"ok": False, "error": "库挂了"} + b = ls._settle_brake() + assert b["high_water"] == 0.0, b # 没写进去就不许当成写进去了 + assert any("高水位" in w for w in b.get("warnings", [])), b + finally: + param_store.set_param, portfolio.positions_view = orig_set, orig_view + + +# ================================================================ +# [E] 安全垫连负天数: 写不上要让日终结算整体报失败 +# ================================================================ +@case("[E1] save_neg_streak 失败 → 返回 ok=False (不再吞成 warning 回 None)") +def _(): + from test_wiring import install_fakes + from app.services import portfolio + from app.repo import pms_repo + install_fakes(prices={}) + orig = pms_repo.set_param + try: + def boom(*a, **kw): + raise RuntimeError("库挂了") + pms_repo.set_param = boom + r = portfolio.save_neg_streak({"600000.SH": 3}) + assert r["ok"] is False and "库挂了" in r["error"], r + finally: + pms_repo.set_param = orig + + +@case("[E2] 连负天数写不上 → daily_settle 的 ok 必须是 False") +def _(): + from test_wiring import install_fakes + from app.services import ledger_service as ls, portfolio + install_fakes(prices={}) + orig = portfolio.save_neg_streak + try: + portfolio.save_neg_streak = lambda m: {"ok": False, "error": "库挂了"} + out = ls.daily_settle() + assert out["ok"] is False, out + assert any("安全垫" in e for e in out["errors"]), out["errors"] + finally: + portfolio.save_neg_streak = orig + + +# ================================================================ +# [F] 一次性守卫计数器: 没写上要留痕 +# ================================================================ +@case("[F1] update_position 影响 0 行 → bump_once_guards 回 ok=False") +def _(): + from test_wiring import install_fakes + from app.services import proposal_service + from app.repo import pms_repo + install_fakes(prices={}) + orig = pms_repo.update_position + try: + pms_repo.update_position = lambda code, **kw: 0 # 持仓行不存在 + r = proposal_service.bump_once_guards("600000.SH", "FILL") + assert r["ok"] is False and "0 行" in r["error"], r + # 不涉及计数器的动作不该被误判成失败 + assert proposal_service.bump_once_guards("600000.SH", "TRIM")["ok"] is True + finally: + pms_repo.update_position = orig + + +@case("[F2] 计数器没写上 → 评审账本留一条 WARN 痕 (指令仍落表, 但纪律失效要有人知道)") +def _(): + from test_wiring import install_fakes + from app.services import proposal_service + from datetime import datetime + fake = install_fakes(prices={"600000.SH": 10.0}) + orig = proposal_service.bump_once_guards + try: + proposal_service.bump_once_guards = ( + lambda code, act, hn=None, now=None: {"ok": False, "fields": {"fill_count": 1}, + "error": "库挂了"}) + iid = proposal_service._make_instruction( + {"ts_code": "600000.SH", "action": "FILL", "side": "buy", "qty": 100, + "reason": "测试", "hard_numbers": {}}, 10.0, datetime.now()) + assert iid, iid + warn = [x for x in fake.ledger if x.get("verdict") == "WARN" + and "一次性守卫" in str(x.get("reason"))] + assert len(warn) == 1, fake.ledger + assert warn[0]["ref_id"] == iid, warn + finally: + proposal_service.bump_once_guards = orig + + +# ================================================================ +# [G] 信号去重键: 落库成功之后才算用掉 +# ================================================================ +@case("[G1] 落指令抛异常 → 当天的去重键不许被烧掉 (下一跳还能重来)") +def _(): + from test_wiring import install_fakes + from app.services import signal_service + from app.core import signal_rules as sr + install_fakes(prices={"600000.SH": 10.0}) + seen, out = set(), {"ignored": 0, "recorded": 0, "exits": [], "proposals": [], + "errors": []} + sig = {"msg_id": "M1", "ts_code": "600000.SH", "action": "SELL", "confidence": 0.95, + "source": "test", "reason": "风控"} + view = {"positions": [{"ts_code": "600000.SH", "total_qty": 1000, "avail_qty": 1000, + "price": 10.0}]} + prm = {} + d = sr.digest(sig, view["positions"][0], prm) + if d["action"] != sr.ACT_EXIT: + return # 规则口径变了就跳过, 不假装测到了 + key = sr.dedup_key(sig, 20260731) + orig = signal_service._make_exit + try: + def boom(*a, **kw): + raise RuntimeError("库挂了") + signal_service._make_exit = boom + try: + signal_service._handle(sig, view, prm, seen, 20260731, False, out) + except RuntimeError: + pass + assert key not in seen, ("落库失败却把去重键用掉了 —— 这条风控卖出信号今天" + "再也不会被消化, 而页面只多一行 error") + finally: + signal_service._make_exit = orig + + +@case("[G2] 落成功后去重键照常生效 (修复不能把去重关掉)") +def _(): + from test_wiring import install_fakes + from app.services import signal_service + from app.core import signal_rules as sr + install_fakes(prices={"600000.SH": 10.0}) + seen, out = set(), {"ignored": 0, "recorded": 0, "exits": [], "proposals": [], + "errors": []} + sig = {"msg_id": "M1", "ts_code": "600000.SH", "action": "SELL", "confidence": 0.95, + "source": "test", "reason": "风控"} + view = {"positions": [{"ts_code": "600000.SH", "total_qty": 1000, "avail_qty": 1000, + "price": 10.0}]} + d = sr.digest(sig, view["positions"][0], {}) + if d["action"] != sr.ACT_EXIT: + return + orig = signal_service._has_inflight + try: + signal_service._has_inflight = lambda c: False + signal_service._handle(sig, view, {}, seen, 20260731, False, out) + assert len(out["exits"]) == 1, out + assert sr.dedup_key(sig, 20260731) in seen, seen + signal_service._handle(sig, view, {}, seen, 20260731, False, out) + assert len(out["exits"]) == 1 and out["ignored"] == 1, out # 第二次被去重挡掉 + finally: + signal_service._has_inflight = orig + + +# ================================================================ +# [H] 日终结算: 对账拒绝时不许报 ok=True +# ================================================================ +@case("[H1] reconcile 回 ok=False → daily_settle 必须 ok=False 并说明原因") +def _(): + from test_wiring import install_fakes + from app.services import ledger_service as ls + install_fakes(prices={}) + orig = ls.reconcile + try: + ls.reconcile = lambda **kw: {"ok": False, "errors": ["两个源都无应答"], + "diffs": [], "fixes": []} + out = ls.daily_settle() + assert out["ok"] is False, out + assert any("对账未完成" in e and "无应答" in e for e in out["errors"]), out["errors"] + finally: + ls.reconcile = orig + + +@case("[H2] 对账连续不一致升到 ERROR → daily_settle 同样不许报成功") +def _(): + from test_wiring import install_fakes + from app.services import ledger_service as ls + install_fakes(prices={}) + orig = ls.reconcile + try: + ls.reconcile = lambda **kw: {"ok": True, "severity": "ERROR", "streak": 3, + "diffs": [{"ts_code": "600000.SH"}], "fixes": []} + out = ls.daily_settle() + assert out["ok"] is False, out + assert any("连续 3 日" in e for e in out["errors"]), out["errors"] + finally: + ls.reconcile = orig + + +# ================================================================ +# [I] 取不到现价的票: 不许拿成本价冒充, 更不许凭空算出安全垫 +# ================================================================ +@case("[I1] 无价的票 cushion_pct 必须是 None, 不能是 0 (0 会被读成「不赚不亏」)") +def _(): + from test_wiring import install_fakes + from app.services import portfolio + fake = install_fakes(prices={"600000.SH": 12.0}) # 000001.SZ 故意没价 + for code, cost in (("600000.SH", 10.0), ("000001.SZ", 10.0)): + fake.insert_lot(ts_code=code, lot_type="BASE", qty=1000, open_price=cost, + open_date="2026-07-01") + fake.update_position(code, total_qty=1000, avail_qty=1000, avg_cost=cost) + v = portfolio.positions_view() + by = {x["ts_code"]: x for x in v["held"]} + assert by["600000.SH"]["price_ok"] is True + assert abs(by["600000.SH"]["cushion_pct"] - 0.2) < 1e-6, by["600000.SH"] + # 取不到价的那只: price 用成本顶着好让市值不塌, 但垫子必须是"不知道" + assert by["000001.SZ"]["price_ok"] is False, by["000001.SZ"] + assert by["000001.SZ"]["cushion_pct"] is None, by["000001.SZ"] + assert "000001.SZ" in v["price_missing"], v.get("price_missing") + + +@case("[I2] 动作引擎跳过无价的票, 且**跳过这件事本身是可见的**") +def _(): + from app.core import action_engine as ae + r = ae.scan(positions=[{"ts_code": "000001.SZ", "total_qty": 1000, "avail_qty": 1000, + "avg_cost": 10.0, "price": 10.0, "price_ok": False, + "cushion_pct": None, "cushion_peak": 0.0}], + params={}, market={}) + assert not r["candidates"], r + assert any(s["ts_code"] == "000001.SZ" and "取不到现价" in s["why"] + for s in r.get("skipped", [])), r + + +@case("[I3] 清仓命令不许悄悄漏掉无价的票, 只数要对得上持仓只数") +def _(): + from app.core import planner as pl + positions = [{"ts_code": "600000.SH", "total_qty": 1000, "avail_qty": 1000, + "price": 10.0, "price_ok": True}, + {"ts_code": "000001.SZ", "total_qty": 500, "avail_qty": 500, + "price": 10.0, "price_ok": False}] + r = pl.plan_liquidate_all(positions=positions, pending_buys=[]) + codes = {i["ts_code"] for i in r["items"] if i["action"] in ("EXIT", "SELL")} + assert codes == {"600000.SH", "000001.SZ"}, r["items"] + need = [i for i in r["items"] if i.get("need_price")] + assert [i["ts_code"] for i in need] == ["000001.SZ"], need + assert r["ok"] is True, r + + +# ================================================================ +# [J] 不追高闸: 拿不到数就说拿不到, 不许当成"通过" +# ================================================================ +@case("[J1] 当日涨幅超上限 → NO_CHASE_DAYUP 拦住 (这道闸得真能拦)") +def _(): + from app.core import rule_gate as rg + r = rg.check(side="buy", action="FILL", qty=100, price=11.0, + ctx={"position": {"total_qty": 0, "avail_qty": 0}, + "caps": None, "params": {"buy_halt_dayup": 0.05}, + "day": {"price": 11.0, "day_chg_from_open": 0.09, "ma5": 11.0}, + "flags": {}}) + assert any("NO_CHASE_DAYUP" in f for f in r["failed"]), r + + +@case("[J2] 取不到当日涨幅 → 留 DAYUP_MISSING 警示, 绝不当成校验通过") +def _(): + from app.core import rule_gate as rg + r = rg.check(side="buy", action="FILL", qty=100, price=11.0, + ctx={"position": {"total_qty": 0, "avail_qty": 0}, + "caps": None, "params": {"buy_halt_dayup": 0.05}, + "day": {"price": 11.0, "day_chg_from_open": None, "ma5": 11.0}, + "flags": {}}) + assert not any("NO_CHASE_DAYUP" in f for f in r["failed"]), r + assert any("DAYUP_MISSING" in w for w in r["warnings"]), r + + +@case("[J3] 自主提议给规则闸的 day 必须是真行情, 不是拿 price 拼出来的空壳") +def _(): + import inspect + from app.services import proposal_service + # 曾经是 {"vwap": price, "day_chg_from_open": None} —— 当日涨幅恒 None, 于是 + # "不追高(涨幅)"这一项对所有自主买入从来没有真正跑过。这里直接验行为: 造一只 + # 当日大涨的票, 走 _route_one, 规则闸必须拿到真涨幅并拦下来。 + from datetime import datetime + from test_wiring import install_fakes + from app.services import market, proposal_service + install_fakes(prices={"600000.SH": 11.0}) + orig = market.day_snapshot + try: + market.day_snapshot = lambda c: {"price": 11.0, "vwap": 10.8, "open": 10.0, + "day_chg_from_open": 0.10, "bars": 60} + mkt = proposal_service._market_ctx([{"ts_code": "600000.SH"}], datetime.now()) + assert mkt["600000.SH"]["day"]["day_chg_from_open"] == 0.10, mkt + # 取快照抛异常也不能让整轮扫描崩, 但要留空让规则闸记 DAYUP_MISSING + def boom(c): + raise RuntimeError("行情库不可用") + market.day_snapshot = boom + mkt = proposal_service._market_ctx([{"ts_code": "600000.SH"}], datetime.now()) + assert mkt["600000.SH"]["day"] == {}, mkt + finally: + market.day_snapshot = orig + + +# ================================================================ +# [K] 参数表读不到时, 安全开关按"拦"而不是按默认值放行 +# ================================================================ +@case("[K1] 参数表读失败 → HALT 开关 fail-closed 取 True (宁可多拦一轮)") +def _(): + from test_wiring import install_fakes + from app.services import param_store + from app.repo import pms_repo + install_fakes(prices={}) + orig, snap = pms_repo.all_params, dict(param_store._cache) + try: + def boom(): + raise RuntimeError("DB 挂了") + pms_repo.all_params = boom + param_store.refresh(force=True) + assert param_store._cache["error"], param_store._cache + assert param_store.get("PMS_GLOBAL_BUY_HALT") is True + assert param_store.get("PMS_GLOBAL_EXEC_HALT") is True + # 非安全开关不受影响, 照常回初值 —— fail-closed 只用在"拦得住"的地方 + assert param_store.get("PMS_AUTONOMY") in ("full", "propose_only", "off") + finally: + pms_repo.all_params = orig + param_store._cache.clear() + param_store._cache.update(snap) + + +@case("[K2] 表读得到、只是没设过这个键 → 照常走默认值 (不能把'没设'当成'读不到')") +def _(): + from test_wiring import install_fakes + from app.services import param_store + install_fakes(prices={}) # 空参数表, 但**读得到** + param_store.refresh(force=True) + assert not param_store._cache["error"], param_store._cache + assert param_store.get("PMS_GLOBAL_BUY_HALT") is not True, "误伤: 没设过被当成读不到" + + +# ================================================================ +def main(): + ok = fail = 0 + for name, fn in CASES: + try: + fn() + print(f" ok {name}") + ok += 1 + except Exception as e: + print(f" FAIL {name}\n {type(e).__name__}: {e}") + fail += 1 + print("-" * 62) + print(f"通过 {ok} 例, 失败 {fail} 例") + if fail: + print("BATCH10 FAIL") + sys.exit(1) + print("BATCH10 PASS") + + +if __name__ == "__main__": + main() diff --git a/scripts/test_batch2_units.py b/scripts/test_batch2_units.py index e70663e..9f202ae 100644 --- a/scripts/test_batch2_units.py +++ b/scripts/test_batch2_units.py @@ -367,10 +367,16 @@ def _(): assert {i["ts_code"] for i in r2["items"]} == {"600000.SH", "000001.SZ"}, r2 assert abs(r2["planned_amount"] - 220_000) < 1e-6 - # 银行占比 220/420 = 52.4% > 40% → 需减 52,000 元 + # 银行占比 220/420 = 52.4% > 40%。 + # **卖出会同时缩小分子和分母**, 所以要解 (sec−x)/(port−x)=cap → x=(sec−cap·port)/(1−cap) + # = 52,000/0.6 = 86,667, 不是 52,000 —— 按 52,000 卖完还是 51.2%, 而命令会报完成 + # (2026-07-31 修; 这条断言原来钉的就是错公式的答案)。 r3 = pl.plan_sector_cap(sector="银行", cap=0.40, positions=ps) assert r3["items"] and all(i["action"] == pl.A_TRIM for i in r3["items"]), r3 - assert abs(r3["target_amount"] - 52_000) < 1e-6, r3 + assert abs(r3["target_amount"] - 86_666.67) < 0.01, r3 + # 减后必须真的落到上限附近 (残留只该来自一手取整) + assert 0.40 <= r3["sector_ratio_after"] <= 0.405, r3["sector_ratio_after"] + assert any("减后仍超上限" in n for n in r3["notes"]), "残留没报出来就等于说完成了" r4 = pl.plan_sector_cap(sector="白酒", cap=0.40, positions=ps) assert r4["items"] == [] and "未超上限" in r4["notes"][0], r4 diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py index 71038c4..d8be8ad 100644 --- a/scripts/test_wiring.py +++ b/scripts/test_wiring.py @@ -898,7 +898,7 @@ def _(): @case("账本服务·对账以下游为准 + 连续不一致升级") def _(): - from app.services import ledger_service as ls + from app.services import ledger_service as ls, param_store from app.repo import downstream_repo fake = install_fakes(prices={"600000.SH": 10.0}) fake.insert_lot(ts_code="600000.SH", lot_type="BASE", qty=1000, open_price=10.0, @@ -915,11 +915,33 @@ def _(): assert any(l["lot_type"] == "RECON" for l in fake.lots) # 修正留痕 assert any(x.get("action") == "RECON" for x in fake.ledger) assert r["severity"] == "WARN" - for _i in range(2): # 连续第 3 日 → ERROR - downstream_repo.fetch_positions = lambda: { - "rows": [{"ts_code": "600000.SH", "qty": 1500 + 100 * (_i + 1)}], - "columns": {"qty": "current_qty"}, "raw_count": 1} + assert r["streak"] == 1, r + + # 每跑一趟对账下游都比账本多 100 股 —— 这样每一趟都真的"不一致", 隔离出 + # 「同一天多趟到底加不加」这一个变量 + nxt = [1500] + + def _drift(): + nxt[0] += 100 + return {"rows": [{"ts_code": "600000.SH", "qty": nxt[0]}], + "columns": {"qty": "current_qty"}, "raw_count": 1} + downstream_repo.fetch_positions = _drift + + # 同一交易日内再对账多少趟, 连续天数都**不动** —— 盘中轻对账每分钟跑一次, + # 按次累加的话"连续 3 日"三分钟就到了, 日报还会写出"连续 175 日"这种数。 + for _ in range(3): r = ls.reconcile() + assert r["diffs"] and r["streak"] == 1 and r["severity"] == "WARN", r + + # 跨交易日才推进: 把"上次推进日"往前拨, 等价于隔了一天再跑。 + # 必须走 set_param (不能直接改 fake.params) —— 一是要过 ParamStore 的缓存失效, + # 二是**顺带守住白名单**: STREAK_YMD 键当初就是漏在白名单外, set_param 静默拒写、 + # prev_ymd 永远读回 0, 按日推进形同虚设。这里 ok 断言就是那道锁。 + for _i in range(2): + w = param_store.set_param(ls.STREAK_YMD_KEY, 20260701 + _i, "test") + assert w["ok"], w + r = ls.reconcile() + assert r["streak"] == 3, r # 连续第 3 日 assert r["severity"] == "ERROR", r finally: downstream_repo.fetch_positions = orig