From 4f2a5455056652fb880ce0b57007d3751715b884 Mon Sep 17 00:00:00 2001 From: zlt Date: Thu, 30 Jul 2026 12:05:22 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E4=BA=8B=E5=AE=9E=E6=BA=90=E5=88=A4?= =?UTF-8?q?=E6=8D=AE:=20=E5=BA=94=E7=AD=94=E7=A9=BA=E9=9B=86=E2=89=A0?= =?UTF-8?q?=E6=B2=A1=E5=BA=94=E7=AD=94,=20force=20=E7=94=9F=E8=B7=AF;=20SR?= =?UTF-8?q?C=5FNONE=20=E4=B8=8B=20force=20=E4=B8=8D=E6=94=BE=E8=A1=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- app/services/ledger_service.py | 36 +++++++++++++++++++++--------- scripts/test_wiring.py | 40 +++++++++++++++++++++++++++------- scripts/ws_smoke.py | 7 ++++-- 3 files changed, 63 insertions(+), 20 deletions(-) diff --git a/app/services/ledger_service.py b/app/services/ledger_service.py index baa39aa..6157511 100644 --- a/app/services/ledger_service.py +++ b/app/services/ledger_service.py @@ -526,8 +526,14 @@ def positions_source() -> dict: f"不同步, 需查明是谁在写那张表"}) return out - if tbl and (tbl.get("rows") or []): - if tbl["columns"].get("qty") is None: + if tbl is not None and not tbl_err: + # **「应答了空」≠「没应答」。** 表查询成功返回 0 行, 是下游给出的一个**有效应答** + # (「我没有持仓」) —— 它可能是真清仓、也可能是它自己读空, 该交给爆炸半径闸按老规矩 + # 处理, 人工确认后 force 能放行。而"两个源都取不到答案"是压根不知道账户状态, 那才 + # 该在对账入口就拒掉。 + # 一开始把两者都归成 source=none, 结果 force 被入口那道拒绝挡死, 07-29 那个 + # 「下游读空绝不清账」用例的 force 分支当场挂了 —— 空集是数据, 不是缺数据。 + if tbl["rows"] and tbl["columns"].get("qty") is None: out["alerts"].append({"level": "ERROR", "code": "TABLE_NO_QTY_COL", "message": "trading_position 未识别出数量列 —— 按 " "QMT_INTERFACE_REQUIREMENTS A1/D1 取 DDL 后把列名" @@ -535,7 +541,7 @@ def positions_source() -> dict: return out # 认不出数量列 = 读不到, 不是没持仓 out.update({"source": SRC_TABLE, "rows": tbl["rows"], "columns": tbl["columns"], "raw_count": tbl.get("raw_count") or len(tbl["rows"])}) - if mode == "ws_first": + if mode == "ws_first" and ws_why: out["alerts"].append({"level": "WARN", "code": "WS_SNAPSHOT_UNAVAILABLE", "message": f"退回 trading_position 表作为事实源: {ws_why}"}) return out @@ -563,17 +569,27 @@ def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict: "[对账·事实源] %s", a["message"]) if src["source"] == SRC_NONE: - # **"什么都读不到" 绝不能当成 "清仓"。** 这里按账本有没有持仓分两级, 不是一律 ERROR: - # 账本也空时 (刚清账、等对端装持仓) 本来就无账可对, 每分钟刷一条 ERROR 只会把真告警 - # 埋掉 —— 与补发期告警限流同一个道理。账本有持仓却读不到事实源, 那才是真要停下来的事。 + # 走到这里意味着**两个源都没有给出应答**(ws 无新鲜快照 + 表查询异常/未启用), 不是 + # "应答了空集"——后者是有效数据, 归 SRC_TABLE 交给爆炸半径闸, 见 positions_source。 + # + # 按账本有没有持仓分两级, 不是一律 ERROR: 账本也空时 (刚清账、等对端装持仓) 本来就 + # 无账可对, 而盘中轻对账每分钟一跳, 刷 ERROR 只会把真告警埋掉 —— 与补发期告警限流 + # 同一个道理。账本有持仓却拿不到任何事实源, 那才是真要停下来的事。 + # + # **这一级连 force 都不放行**, 与爆炸半径闸不同。force 的语义是「人工已确认下游读数 + # 正确」, 而这里根本没有读数可供确认 —— 没有任何数字, 人也无从确认。真要清账走 + # scripts/reset_ledger.py 那条明路, 别拿一个空壳子当"下游事实"去核销批次。 held = [p for p in pms_repo.list_positions() if int(p.get("total_qty") or 0) > 0] if held: out.update({"severity": rc.SEV_ERROR, "fixes": [], - "blocked": {"why": "拿不到任何持仓事实源, 而本端有持仓 —— " - "拒绝对账 (读不到 ≠ 清仓)", "held": len(held)}}) - logger.error("[对账] 拒绝对账: 无事实源而本端有 %s 只持仓", len(held)) + "blocked": {"why": "两个事实源都没有应答 (不是应答了空集), 而本端有" + "持仓 —— 拒绝对账。读不到 ≠ 清仓; force 在此不放行, " + "要清账用 scripts/reset_ledger.py", + "held": len(held), "force_ignored": bool(force)}}) + logger.error("[对账] 拒绝对账: 无事实源应答而本端有 %s 只持仓 (force=%s 不放行)", + len(held), force) else: - out["note"] = "事实源与账本都空, 无可对之账 (等对端装持仓)" + out["note"] = "两个事实源都没有应答, 账本也空 —— 无可对之账 (等对端装持仓)" return out # 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、 diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py index 655e6be..1f4fd25 100644 --- a/scripts/test_wiring.py +++ b/scripts/test_wiring.py @@ -1193,27 +1193,51 @@ def _(): assert "过期" in msg, src["alerts"] -@case("对账事实源·两个源都拿不到 + 本端有持仓 → **拒绝对账**, 一股都不许核销") +@case("对账事实源·表应答空集仍算有效应答 → 归 table (空集是数据, 不是缺数据)") def _(): + from app.services import ledger_service as ls + install_fakes(prices={"600000.SH": 10.0}) # 默认桩: fetch_positions 成功返回 rows=[] + src = ls.positions_source() + # 这一条守着 force 的生路: 归 table 才会走到爆炸半径闸, 人工确认后 force 能放行; + # 若把"应答了空"误判成 source=none, 对账入口就把 force 挡死了 (曾经真挡死过) + assert src["source"] == ls.SRC_TABLE, src + assert src["rows"] == [], src + + +@case("对账事实源·两个源都**无应答** + 本端有持仓 → 入口就拒, **force 也不放行**") +def _(): + from app.repo import downstream_repo from app.services import ledger_service as ls fake = install_fakes(prices={"600000.SH": 10.0}) fake.insert_lot(ts_code="600000.SH", lot_type="BASE", qty=1000, open_price=10.0, open_date="2026-07-01") fake.update_position("600000.SH", total_qty=1000, avail_qty=1000) - src = ls.positions_source() # 没塞快照, 表也空 + + def boom(): + raise RuntimeError("Can't connect to MySQL server on '192.168.16.153'") + downstream_repo.fetch_positions = boom # 没塞快照 + 表查询异常 = 谁都没应答 + src = ls.positions_source() assert src["source"] == ls.SRC_NONE, src - r = ls.reconcile(apply_fix=True) - assert r.get("blocked"), "读不到 ≠ 清仓, 必须拒绝对账" - assert r["severity"] == "ERROR" and not r["fixes"], r - assert fake.positions["600000.SH"]["total_qty"] == 1000, "持仓不得被动过" + # force 的语义是「人工已确认下游读数正确」, 而这里根本没有读数可供确认 —— 没有任何数字, + # 人也无从确认。真要清账走 reset_ledger.py, 别拿空壳子当"下游事实"去核销批次。 + for force in (False, True): + r = ls.reconcile(apply_fix=True, force=force) + assert r.get("blocked"), f"force={force} 也必须拒绝对账" + assert r["severity"] == "ERROR" and not r["fixes"], r + assert fake.positions["600000.SH"]["total_qty"] == 1000, "持仓一股都不许被动" -@case("对账事实源·两个源都空 + 账本也空 → 不报 ERROR (刚清账等对端装持仓)") +@case("对账事实源·两个源都无应答 + 账本也空 → 不报 ERROR (刚清账等对端装持仓)") def _(): + from app.repo import downstream_repo from app.services import ledger_service as ls install_fakes() + + def boom(): + raise RuntimeError("库不可达") + downstream_repo.fetch_positions = boom r = ls.reconcile(apply_fix=True) - # 每分钟一跳的轻对账在这个状态下会反复走到这儿, 刷 ERROR 只会把真告警埋掉 + # 盘中轻对账每分钟一跳, 这个状态下反复刷 ERROR 只会把真告警埋掉 assert not r.get("blocked") and r["severity"] == "OK", r assert "无可对之账" in (r.get("note") or ""), r diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index 803691c..84efe8a 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -327,7 +327,9 @@ def cmd_inbox(args): for r in rows: mark = {0: "待入账", 1: "已入账", 2: "已消化", 3: "挂起"}.get(int(r.get("processed") or 0), "?") body = json.dumps(r.get("payload") or {}, ensure_ascii=False) - print(f" seq={r['seq']:<8} {r['msg_type']:<15} {mark} {body[:110]}") + # 默认 110 字够看 trade/order_update; snapshot 的 items 一截就把 cost_price 切掉了, + # 而那正是接管持仓时最要紧的一个字段 —— 用 --width 放宽 + print(f" seq={r['seq']:<8} {r['msg_type']:<15} {mark} {body[:int(args.width)]}") return 0 @@ -365,8 +367,9 @@ def main(): p = sub.add_parser("inbox", help="最近上行消息") p.add_argument("--limit", default=20) p.add_argument("--type", default=None, - help="只看某一类, 如 trade / ack / order_update / reject " + help="只看某一类, 如 trade / ack / order_update / reject / snapshot " "(不填会被 pong 刷屏)") + p.add_argument("--width", default=110, help="payload 打印宽度 (看 snapshot 建议 600)") p.set_defaults(fn=cmd_inbox) p = sub.add_parser("rewind", help="回退 seq 水位, 逼对端补发 (协议 §6.1; 须先停 pms-ws)")