修事实源判据: 应答空集≠没应答, force 生路; SRC_NONE 下 force 不放行
This commit is contained in:
parent
e32b6df05b
commit
4f2a545505
|
|
@ -526,8 +526,14 @@ def positions_source() -> dict:
|
||||||
f"不同步, 需查明是谁在写那张表"})
|
f"不同步, 需查明是谁在写那张表"})
|
||||||
return out
|
return out
|
||||||
|
|
||||||
if tbl and (tbl.get("rows") or []):
|
if tbl is not None and not tbl_err:
|
||||||
if tbl["columns"].get("qty") is None:
|
# **「应答了空」≠「没应答」。** 表查询成功返回 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",
|
out["alerts"].append({"level": "ERROR", "code": "TABLE_NO_QTY_COL",
|
||||||
"message": "trading_position 未识别出数量列 —— 按 "
|
"message": "trading_position 未识别出数量列 —— 按 "
|
||||||
"QMT_INTERFACE_REQUIREMENTS A1/D1 取 DDL 后把列名"
|
"QMT_INTERFACE_REQUIREMENTS A1/D1 取 DDL 后把列名"
|
||||||
|
|
@ -535,7 +541,7 @@ def positions_source() -> dict:
|
||||||
return out # 认不出数量列 = 读不到, 不是没持仓
|
return out # 认不出数量列 = 读不到, 不是没持仓
|
||||||
out.update({"source": SRC_TABLE, "rows": tbl["rows"], "columns": tbl["columns"],
|
out.update({"source": SRC_TABLE, "rows": tbl["rows"], "columns": tbl["columns"],
|
||||||
"raw_count": tbl.get("raw_count") or len(tbl["rows"])})
|
"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",
|
out["alerts"].append({"level": "WARN", "code": "WS_SNAPSHOT_UNAVAILABLE",
|
||||||
"message": f"退回 trading_position 表作为事实源: {ws_why}"})
|
"message": f"退回 trading_position 表作为事实源: {ws_why}"})
|
||||||
return out
|
return out
|
||||||
|
|
@ -563,17 +569,27 @@ def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict:
|
||||||
"[对账·事实源] %s", a["message"])
|
"[对账·事实源] %s", a["message"])
|
||||||
|
|
||||||
if src["source"] == SRC_NONE:
|
if src["source"] == SRC_NONE:
|
||||||
# **"什么都读不到" 绝不能当成 "清仓"。** 这里按账本有没有持仓分两级, 不是一律 ERROR:
|
# 走到这里意味着**两个源都没有给出应答**(ws 无新鲜快照 + 表查询异常/未启用), 不是
|
||||||
# 账本也空时 (刚清账、等对端装持仓) 本来就无账可对, 每分钟刷一条 ERROR 只会把真告警
|
# "应答了空集"——后者是有效数据, 归 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]
|
held = [p for p in pms_repo.list_positions() if int(p.get("total_qty") or 0) > 0]
|
||||||
if held:
|
if held:
|
||||||
out.update({"severity": rc.SEV_ERROR, "fixes": [],
|
out.update({"severity": rc.SEV_ERROR, "fixes": [],
|
||||||
"blocked": {"why": "拿不到任何持仓事实源, 而本端有持仓 —— "
|
"blocked": {"why": "两个事实源都没有应答 (不是应答了空集), 而本端有"
|
||||||
"拒绝对账 (读不到 ≠ 清仓)", "held": len(held)}})
|
"持仓 —— 拒绝对账。读不到 ≠ 清仓; force 在此不放行, "
|
||||||
logger.error("[对账] 拒绝对账: 无事实源而本端有 %s 只持仓", len(held))
|
"要清账用 scripts/reset_ledger.py",
|
||||||
|
"held": len(held), "force_ignored": bool(force)}})
|
||||||
|
logger.error("[对账] 拒绝对账: 无事实源应答而本端有 %s 只持仓 (force=%s 不放行)",
|
||||||
|
len(held), force)
|
||||||
else:
|
else:
|
||||||
out["note"] = "事实源与账本都空, 无可对之账 (等对端装持仓)"
|
out["note"] = "两个事实源都没有应答, 账本也空 —— 无可对之账 (等对端装持仓)"
|
||||||
return out
|
return out
|
||||||
|
|
||||||
# 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、
|
# 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、
|
||||||
|
|
|
||||||
|
|
@ -1193,27 +1193,51 @@ def _():
|
||||||
assert "过期" in msg, src["alerts"]
|
assert "过期" in msg, src["alerts"]
|
||||||
|
|
||||||
|
|
||||||
@case("对账事实源·两个源都拿不到 + 本端有持仓 → **拒绝对账**, 一股都不许核销")
|
@case("对账事实源·表应答空集仍算有效应答 → 归 table (空集是数据, 不是缺数据)")
|
||||||
def _():
|
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
|
from app.services import ledger_service as ls
|
||||||
fake = install_fakes(prices={"600000.SH": 10.0})
|
fake = install_fakes(prices={"600000.SH": 10.0})
|
||||||
fake.insert_lot(ts_code="600000.SH", lot_type="BASE", qty=1000, open_price=10.0,
|
fake.insert_lot(ts_code="600000.SH", lot_type="BASE", qty=1000, open_price=10.0,
|
||||||
open_date="2026-07-01")
|
open_date="2026-07-01")
|
||||||
fake.update_position("600000.SH", total_qty=1000, avail_qty=1000)
|
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
|
assert src["source"] == ls.SRC_NONE, src
|
||||||
r = ls.reconcile(apply_fix=True)
|
# force 的语义是「人工已确认下游读数正确」, 而这里根本没有读数可供确认 —— 没有任何数字,
|
||||||
assert r.get("blocked"), "读不到 ≠ 清仓, 必须拒绝对账"
|
# 人也无从确认。真要清账走 reset_ledger.py, 别拿空壳子当"下游事实"去核销批次。
|
||||||
assert r["severity"] == "ERROR" and not r["fixes"], r
|
for force in (False, True):
|
||||||
assert fake.positions["600000.SH"]["total_qty"] == 1000, "持仓不得被动过"
|
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 _():
|
def _():
|
||||||
|
from app.repo import downstream_repo
|
||||||
from app.services import ledger_service as ls
|
from app.services import ledger_service as ls
|
||||||
install_fakes()
|
install_fakes()
|
||||||
|
|
||||||
|
def boom():
|
||||||
|
raise RuntimeError("库不可达")
|
||||||
|
downstream_repo.fetch_positions = boom
|
||||||
r = ls.reconcile(apply_fix=True)
|
r = ls.reconcile(apply_fix=True)
|
||||||
# 每分钟一跳的轻对账在这个状态下会反复走到这儿, 刷 ERROR 只会把真告警埋掉
|
# 盘中轻对账每分钟一跳, 这个状态下反复刷 ERROR 只会把真告警埋掉
|
||||||
assert not r.get("blocked") and r["severity"] == "OK", r
|
assert not r.get("blocked") and r["severity"] == "OK", r
|
||||||
assert "无可对之账" in (r.get("note") or ""), r
|
assert "无可对之账" in (r.get("note") or ""), r
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -327,7 +327,9 @@ def cmd_inbox(args):
|
||||||
for r in rows:
|
for r in rows:
|
||||||
mark = {0: "待入账", 1: "已入账", 2: "已消化", 3: "挂起"}.get(int(r.get("processed") or 0), "?")
|
mark = {0: "待入账", 1: "已入账", 2: "已消化", 3: "挂起"}.get(int(r.get("processed") or 0), "?")
|
||||||
body = json.dumps(r.get("payload") or {}, ensure_ascii=False)
|
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
|
return 0
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -365,8 +367,9 @@ def main():
|
||||||
p = sub.add_parser("inbox", help="最近上行消息")
|
p = sub.add_parser("inbox", help="最近上行消息")
|
||||||
p.add_argument("--limit", default=20)
|
p.add_argument("--limit", default=20)
|
||||||
p.add_argument("--type", default=None,
|
p.add_argument("--type", default=None,
|
||||||
help="只看某一类, 如 trade / ack / order_update / reject "
|
help="只看某一类, 如 trade / ack / order_update / reject / snapshot "
|
||||||
"(不填会被 pong 刷屏)")
|
"(不填会被 pong 刷屏)")
|
||||||
|
p.add_argument("--width", default=110, help="payload 打印宽度 (看 snapshot 建议 600)")
|
||||||
p.set_defaults(fn=cmd_inbox)
|
p.set_defaults(fn=cmd_inbox)
|
||||||
|
|
||||||
p = sub.add_parser("rewind", help="回退 seq 水位, 逼对端补发 (协议 §6.1; 须先停 pms-ws)")
|
p = sub.add_parser("rewind", help="回退 seq 水位, 逼对端补发 (协议 §6.1; 须先停 pms-ws)")
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue