补协议§6.2快照对账: query轮询+ws为主表兜底的事实源仲裁

This commit is contained in:
zlt 2026-07-30 11:59:12 +08:00
parent a25c4d1540
commit e32b6df05b
8 changed files with 485 additions and 14 deletions

View File

@ -425,6 +425,106 @@ def dedup_key(type_: str, payload: dict):
return None return None
# ---- snapshot(kind=positions) 解析 (§5.7) ----------------------------------
# 字段名给候选而不是写死, 沿用 downstream_repo 那套探测思路 (首项为协议正文的名字)。
# 理由很实际: 协议正文写的是 ts_code/total_qty/avail_qty, 但对端的实现是从券商结构体
# 映射过来的, 极可能带着 stock_code/total_quantity 这类原生名。字段名不合就整批读不到,
# 而"读不到"和"真的没持仓"在对账那头长成同一个样子 —— 这正是最不能混的两件事。
_SNAP_LIST_KEYS = ("items", "positions", "rows", "data", "list")
_SNAP_CODE = ("ts_code", "stock_code", "code", "security_code", "symbol")
_SNAP_QTY = ("total_qty", "total_quantity", "qty", "quantity", "volume", "current_qty",
"position_qty", "hold_qty")
_SNAP_AVAIL = ("avail_qty", "available_qty", "available_quantity", "can_use_volume",
"sellable_qty", "enable_amount")
_SNAP_FROZEN = ("frozen_qty", "frozen_quantity", "frozen_volume", "freeze_qty")
_SNAP_COST = ("cost_price", "avg_cost", "cost", "open_price", "position_cost")
_SNAP_PRICE = ("market_price", "last_price", "current_price", "price")
def _first(d: dict, names, cast, default=None):
low = {str(k).lower(): v for k, v in (d or {}).items()}
for n in names:
if n in low and low[n] not in (None, ""):
try:
return cast(low[n]), n
except (TypeError, ValueError):
return default, n
return default, None
def parse_positions_snapshot(payload: dict) -> dict:
"""`snapshot{kind:"positions"}` → 与 downstream_repo.fetch_positions() **同形**的结构。
同形是刻意的: 对账引擎不该关心持仓是从 ws 快照来的还是从 trading_position 表来的,
换源只换取数, 不动比对与修正逻辑
返回 {"rows": [{ts_code, qty, avail_qty, cost, frozen, price}], "columns": {...},
"raw_count": n, "as_of": epoch_ms, "unknown_fields": [...]}
`columns` 里回填**实际命中的字段名**, 与表通道的语义一致 页面导出下游表结构
对账的未识别出数量列告警都靠它命中不到 qty `columns["qty"] is None`,
上层会据此判定读不到而不是没持仓
**代码只做大小写归一, 不转点式** 协议 §8 规定 ts_code 就是点式, 但对端万一给了前缀式
(`SH600000`), 转换要用 downstream_repo.to_dot() 本模块是零外部依赖的纯逻辑, 不引
repo , 所以归一到 PMS 内部口径由调用方做 ( ledger_service.positions_source)
自带一份 to_dot 的话两处实现迟早漂开, 那种 bug 只会在某个冷门板块上冒出来
"""
pl = payload or {}
items = None
for k in _SNAP_LIST_KEYS:
v = pl.get(k)
if isinstance(v, list):
items = v
break
if items is None:
return {"rows": [], "columns": {"qty": None, "avail": None, "cost": None},
"raw_count": 0, "as_of": int(pl.get("as_of") or 0),
"unknown_fields": sorted(str(k) for k in pl.keys())}
rows, cols = [], {"code": None, "qty": None, "avail": None, "cost": None,
"frozen": None, "price": None}
for it in items:
if not isinstance(it, dict):
continue
code, c_code = _first(it, _SNAP_CODE, str, "")
qty, c_qty = _first(it, _SNAP_QTY, lambda v: int(float(v)), None)
avail, c_av = _first(it, _SNAP_AVAIL, lambda v: int(float(v)), None)
frozen, c_fz = _first(it, _SNAP_FROZEN, lambda v: int(float(v)), None)
cost, c_cs = _first(it, _SNAP_COST, float, None)
price, c_px = _first(it, _SNAP_PRICE, float, None)
for key, hit in (("code", c_code), ("qty", c_qty), ("avail", c_av),
("cost", c_cs), ("frozen", c_fz), ("price", c_px)):
if hit and not cols[key]:
cols[key] = hit
code = (code or "").strip().upper()
if not code:
continue
rows.append({"ts_code": code, "qty": qty, "avail_qty": avail, "cost": cost,
"frozen": frozen, "price": price})
return {"rows": rows, "columns": cols, "raw_count": len(items),
"as_of": int(pl.get("as_of") or 0), "unknown_fields": []}
def parse_funds_snapshot(payload: dict) -> dict:
"""`snapshot{kind:"funds"}` → 归一后的资金字段 (§5.7)。
暂无消费方: `calibrate_fees` 的公式要总资产变动 成交净额, 需要**两天**的快照差
再扣掉入出金, 不是接上字段就能算对先解析落库等口径定了再接, 别硬凑一个看起来
有结果的数 费用本来就不进成本, 晚校准无风险 (协议 §5.5)
"""
d = (payload or {}).get("data")
d = d if isinstance(d, dict) else (payload or {})
out = {"as_of": int((payload or {}).get("as_of") or 0)}
for key, names in (("total_asset", ("total_asset", "total_assets", "asset")),
("available_cash", ("available_cash", "avail_cash", "cash")),
("frozen_cash", ("frozen_cash", "freeze_cash")),
("sell_return_today", ("sell_return_today", "sell_return")),
("market_value", ("market_value", "mv"))):
out[key], _ = _first(d, names, float, None)
return out
def trade_amount_ok(payload: dict, tol: float = 0.01) -> bool: def trade_amount_ok(payload: dict, tol: float = 0.01) -> bool:
"""§5.5: amount 应等于 price × qty, 差异 > 0.01 元告警。""" """§5.5: amount 应等于 price × qty, 差异 > 0.01 元告警。"""
try: try:

View File

@ -281,6 +281,31 @@ def inbox_orphan_count() -> int:
return int((r or {}).get("n") or 0) return int((r or {}).get("n") or 0)
def latest_snapshot(kind: str, *, scan: int = 60) -> dict:
"""取最近一份 `snapshot{kind:...}` 上行 (协议 §5.7)。返回 None 表示没有。
返回 {"seq", "payload", "received_at", "age_sec"} `age_sec` **本端** received_at
, 不用 payload 里的 as_of: 新鲜度是本端的安全判据, 拿对端时钟算会把时钟漂移混进来
(实测两机差约 0.9 , 真漂起来就不止了)as_of 留在 payload 里供展示与对端自证
positions / funds / orders 三种 kind 混在同一个 msg_type , 所以只能先按 seq 倒序捞
一批再在 Python 侧筛 单表查询纪律不允许把 JSON 里的 kind 写进 WHERE
"""
rows = fetch_all("SELECT seq, payload_json, received_at FROM pms_qmt_inbox "
"WHERE msg_type = 'snapshot' ORDER BY seq DESC LIMIT :n",
{"n": int(scan)})
want = str(kind or "").lower()
for r in rows:
pl = _loads(r.get("payload_json"), {}) or {}
if str(pl.get("kind") or "").lower() != want:
continue
rec = r.get("received_at")
age = (_NOW() - rec).total_seconds() if rec else None
return {"seq": int(r["seq"]), "payload": pl, "received_at": rec,
"age_sec": round(age, 1) if age is not None else None}
return None
def inbox_stats_above(seq: int) -> list: def inbox_stats_above(seq: int) -> list:
"""seq 之上的上行消息按 类型×入账状态 汇总。给 rewind 判"删了会不会出事"用。""" """seq 之上的上行消息按 类型×入账状态 汇总。给 rewind 判"删了会不会出事"用。"""
return fetch_all( return fetch_all(

View File

@ -22,6 +22,7 @@ from app.core import command_spec as cs
from app.core import cushion as cu from app.core import cushion as cu
from app.core import recon as rc from app.core import recon as rc
from app.core import tradedays as td from app.core import tradedays as td
from app.core import ws_codec as wsc
from app.repo import downstream_repo, pms_repo, qmt_repo from app.repo import downstream_repo, pms_repo, qmt_repo
from app.services import market, param_store, portfolio from app.services import market, param_store, portfolio
@ -447,23 +448,132 @@ def _recon_blast_guard(book: list, ds: dict, diffs: list, *, force: bool = False
return None return None
SRC_WS, SRC_TABLE, SRC_NONE = "ws", "table", "none"
def positions_source() -> dict:
"""对账的持仓事实源:**ws 快照为主、`trading_position` 表为兜底**2026-07-30 拍板)。
为什么要有这一层
----------------
协议 §6.2 定的对账事实源是 ws `query_positions` `snapshot{kind:"positions"}`;
而目标架构里 `trading_service` 全量退出业务只留看板, **那张表在新架构下没有明确的
写入方**07-30 实测: QMT 已切到模拟仓且功能正常, `trading_position` 却是空的
(`fetch_positions` columns 三个 None 就是"表里一行都没有"的铁证)
只认表的话, 账本永远建不起来
三条仲裁规则**两个源不一致时不许静默挑一个** 那会变成"两个同名不同物":
1. ws 快照新鲜 ws若表也非空且与 ws 对不上, **照样用 ws, 但记一条告警**
列出差异只数 那说明表的写入方与 QMT 已经不同步, 是要修的事, 不是噪音
2. ws 快照缺失/过期/字段不认 退回表, 并说明退回的原因 (三种原因处理起来完全不同:
没接通要找对端, 过期要看 pms-ws 活没活, 字段不认要补 ws_codec 的候选名)
3. 两个源都拿不到 `source=none`**这不等于"清仓"**, 上层必须据此拒绝改账
新鲜度按本端 `received_at` , 不用 payload `as_of` 理由见 qmt_repo.latest_snapshot
"""
mode = param_store.get("PMS_RECON_SOURCE", "ws_first") or "ws_first"
max_age = param_store.get_int("PMS_RECON_WS_SNAPSHOT_MAX_AGE_SEC", 900)
out = {"source": SRC_NONE, "rows": [], "columns": {"qty": None, "avail": None, "cost": None},
"raw_count": 0, "as_of": 0, "age_sec": None, "alerts": [], "mode": mode}
ws_snap, ws_why = None, ""
if mode in ("ws_first", "ws_only"):
try:
snap = qmt_repo.latest_snapshot("positions")
except Exception as e:
snap, ws_why = None, f"读 ws 快照失败: {type(e).__name__}: {e}"
if not snap:
ws_why = ws_why or ("ws 从未回过 positions 快照 —— 确认 pms-ws 在跑、"
"PMS_QMT_WS_ENABLED 已开, 且对端实现了 query_positions")
elif snap.get("age_sec") is not None and snap["age_sec"] > max_age:
ws_why = (f"ws 快照已过期 ({snap['age_sec']:.0f}s > {max_age}s) —— "
f"pms-ws 可能没在跑, 或对端不再回应 query_positions")
else:
parsed = wsc.parse_positions_snapshot(snap["payload"])
if parsed["raw_count"] and parsed["columns"].get("qty") is None:
ws_why = (f"ws 快照有 {parsed['raw_count']} 个条目却认不出数量列 —— "
f"把对端实际字段名补进 ws_codec._SNAP_QTY")
else:
# 归一到 PMS 内部点式口径。协议 §8 规定就是点式, 但对端给前缀式 (SH600000)
# 时若不转, diff 会拿 "SH600000" 去比账本里的 "600000.SH" —— 结果是**每一只
# 都对不上**: 账本那只判"下游没有了"要核销, ws 那只判"新持仓"要补。
# 一次代码格式不一致就能造出一轮双向全量重写, 比读空还狠。
for r in parsed["rows"]:
r["ts_code"] = downstream_repo.to_dot(r["ts_code"])
ws_snap = {**parsed, "age_sec": snap.get("age_sec"), "seq": snap.get("seq")}
tbl, tbl_err = None, ""
if mode in ("ws_first", "table_only"):
try:
tbl = downstream_repo.fetch_positions()
except Exception as e:
tbl_err = f"读 trading_position 失败: {type(e).__name__}: {e}"
if ws_snap:
out.update({"source": SRC_WS, "rows": ws_snap["rows"], "columns": ws_snap["columns"],
"raw_count": ws_snap["raw_count"], "as_of": ws_snap["as_of"],
"age_sec": ws_snap["age_sec"], "seq": ws_snap.get("seq")})
if tbl and (tbl.get("rows") or []):
ws_map = {r["ts_code"]: int(r.get("qty") or 0) for r in ws_snap["rows"]}
tb_map = {r["ts_code"]: int(r.get("qty") or 0) for r in tbl["rows"]}
gap = [c for c in set(ws_map) | set(tb_map) if ws_map.get(c, 0) != tb_map.get(c, 0)]
if gap:
out["alerts"].append(
{"level": "WARN", "code": "SOURCE_DISAGREE",
"message": f"ws 快照与 trading_position 对不上 {len(gap)}"
f"(ws {len(ws_map)} 只 / 表 {len(tb_map)} 只): "
f"{sorted(gap)[:8]}。**已按 ws 为准**; 表的写入方与 QMT "
f"不同步, 需查明是谁在写那张表"})
return out
if tbl and (tbl.get("rows") or []):
if 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 后把列名"
"补进 downstream_repo.QTY_CANDIDATES"})
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":
out["alerts"].append({"level": "WARN", "code": "WS_SNAPSHOT_UNAVAILABLE",
"message": f"退回 trading_position 表作为事实源: {ws_why}"})
return out
out["alerts"].append({"level": "WARN", "code": "NO_POSITION_SOURCE",
"message": f"两个事实源都拿不到持仓 —— ws: {ws_why or '未启用'}; "
f"表: {tbl_err or '空集'}。**这不等于清仓**"})
return out
def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict: def reconcile(*, apply_fix: bool = True, force: bool = False) -> dict:
"""账本 vs 下游持仓, 以下游为准修正并留痕。 """账本 vs 下游持仓, 以下游为准修正并留痕。事实源见 positions_source()。
force=True 绕过爆炸半径限制 ( _recon_blast_guard) 只在人工确认下游读数确实 force=True 绕过爆炸半径限制 ( _recon_blast_guard) 只在人工确认下游读数确实
正确之后使用 正确之后使用
""" """
out = {"ok": True, "diffs": [], "fixes": [], "errors": [], "columns": {}, "severity": rc.SEV_OK} out = {"ok": True, "diffs": [], "fixes": [], "errors": [], "columns": {}, "severity": rc.SEV_OK}
try: src = positions_source()
ds = downstream_repo.fetch_positions() ds = {"rows": src["rows"], "columns": src["columns"], "raw_count": src["raw_count"]}
except Exception as e: out.update({"source": src["source"], "source_mode": src["mode"],
out.update({"ok": False, "errors": [f"读 trading_position 失败: {type(e).__name__}: {e}"]}) "source_age_sec": src.get("age_sec"), "source_alerts": src["alerts"],
return out "columns": src["columns"]})
out["columns"] = ds["columns"] for a in src["alerts"]:
if ds["rows"] and ds["columns"].get("qty") is None: (logger.error if a["level"] == "ERROR" else logger.warning)(
out.update({"ok": False, "errors": [ "[对账·事实源] %s", a["message"])
"下游持仓表未识别出数量列 —— 请按 QMT_INTERFACE_REQUIREMENTS A1/D1 取得 DDL 后, "
"把列名补进 downstream_repo.QTY_CANDIDATES"]}) if src["source"] == SRC_NONE:
# **"什么都读不到" 绝不能当成 "清仓"。** 这里按账本有没有持仓分两级, 不是一律 ERROR:
# 账本也空时 (刚清账、等对端装持仓) 本来就无账可对, 每分钟刷一条 ERROR 只会把真告警
# 埋掉 —— 与补发期告警限流同一个道理。账本有持仓却读不到事实源, 那才是真要停下来的事。
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))
else:
out["note"] = "事实源与账本都空, 无可对之账 (等对端装持仓)"
return out return out
# 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、 # 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、

View File

@ -106,6 +106,7 @@ class WsRunner:
"beat_sec": param_store.get_int("PMS_QMT_HEARTBEAT_DB_SEC", 2), "beat_sec": param_store.get_int("PMS_QMT_HEARTBEAT_DB_SEC", 2),
"max_attempts": param_store.get_int("PMS_QMT_SEND_MAX_ATTEMPTS", 3), "max_attempts": param_store.get_int("PMS_QMT_SEND_MAX_ATTEMPTS", 3),
"connect_timeout": param_store.get_int("PMS_QMT_CONNECT_TIMEOUT_SEC", 10), "connect_timeout": param_store.get_int("PMS_QMT_CONNECT_TIMEOUT_SEC", 10),
"query_interval_sec": param_store.get_int("PMS_QMT_QUERY_INTERVAL_SEC", 300),
} }
try: try:
self._params = await _db(_load) self._params = await _db(_load)
@ -115,7 +116,8 @@ class WsRunner:
"url": settings.PMS_QMT_WS_URL, "heartbeat_sec": 5, "url": settings.PMS_QMT_WS_URL, "heartbeat_sec": 5,
"idle_timeout_sec": 15, "ack_batch": 20, "idle_timeout_sec": 15, "ack_batch": 20,
"ack_interval_sec": 2.0, "outbox_poll_sec": 0.5, "ack_interval_sec": 2.0, "outbox_poll_sec": 0.5,
"beat_sec": 2, "max_attempts": 3, "connect_timeout": 10} "beat_sec": 2, "max_attempts": 3, "connect_timeout": 10,
"query_interval_sec": 300}
logger.warning("参数刷新失败, 沿用上一份: %s", _brief_err(e)) logger.warning("参数刷新失败, 沿用上一份: %s", _brief_err(e))
@staticmethod @staticmethod
@ -287,7 +289,8 @@ class WsRunner:
tasks = [asyncio.create_task(self._reader_loop(peer), name="reader"), tasks = [asyncio.create_task(self._reader_loop(peer), name="reader"),
asyncio.create_task(self._pinger_loop(seed), name="pinger"), asyncio.create_task(self._pinger_loop(seed), name="pinger"),
asyncio.create_task(self._outbox_loop(seed), name="outbox"), asyncio.create_task(self._outbox_loop(seed), name="outbox"),
asyncio.create_task(self._ack_loop(seed), name="ack")] asyncio.create_task(self._ack_loop(seed), name="ack"),
asyncio.create_task(self._query_loop(seed), name="query")]
# 停机信号单独占一个 waiter: reader 阻塞在 15 秒 recv 上, 光等它自己醒来会让 # 停机信号单独占一个 waiter: reader 阻塞在 15 秒 recv 上, 光等它自己醒来会让
# 优雅退出白白拖满一个 idle 周期。谁先结束就收工, 剩下的直接 cancel。 # 优雅退出白白拖满一个 idle 周期。谁先结束就收工, 剩下的直接 cancel。
stopper = asyncio.create_task(self._stop.wait(), name="stopper") stopper = asyncio.create_task(self._stop.wait(), name="stopper")
@ -587,6 +590,21 @@ class WsRunner:
pl.get("trade_no"), json.dumps(pl, ensure_ascii=False)[:200]) pl.get("trade_no"), json.dumps(pl, ensure_ascii=False)[:200])
logger.info("[trade] %s %s %s股 @%s (待 worker 入账)", iid, logger.info("[trade] %s %s %s股 @%s (待 worker 入账)", iid,
pl.get("ts_code"), pl.get("qty"), pl.get("price")) pl.get("ts_code"), pl.get("qty"), pl.get("price"))
elif type_ == wsc.T_SNAPSHOT:
# 只记「拿到了、几条、多新」, 不在这里解析入账 —— 快照的消费方是 worker 里的
# reconcile (从 pms_qmt_inbox 读最近一份), 与出口队列同一个道理: ws 进程只做
# 通道, 账本变更保持单一来源, 也不至于让一次比对把心跳拖到超时断连。
kind = str(pl.get("kind") or "?")
n = len(pl.get("items") or []) if kind == "positions" else None
self._stat[f"snap_{kind}_at"] = datetime.now().strftime("%H:%M:%S")
if n is not None:
self._stat["snap_positions_n"] = n
logger.info("[snapshot] kind=%s%s as_of=%s", kind,
f" 条目={n}" if n is not None else "", pl.get("as_of"))
elif type_ in (wsc.T_POS_UPDATE, wsc.T_FUNDS_UPDATE):
# 增量推送: 只当"有变动了"的提示。**不拿它攒对账视图**, 理由见 _query_loop。
self._stat[f"{type_}_at"] = datetime.now().strftime("%H:%M:%S")
self._stat[f"{type_}_n"] = int(self._stat.get(f"{type_}_n") or 0) + 1
elif type_ == wsc.T_PERSIST: elif type_ == wsc.T_PERSIST:
if not pl.get("ok"): if not pl.get("ok"):
# §5.6: 落库失败不影响成交事实, 但要留一条下游不一致告警 # §5.6: 落库失败不影响成交事实, 但要留一条下游不一致告警
@ -677,6 +695,41 @@ class WsRunner:
return return
await self._send(wsc.T_PING, {}, seed) await self._send(wsc.T_PING, {}, seed)
async def _query_loop(self, seed: str):
"""周期拉**全量**快照 (协议 §6.2「快照对账」): query_positions + query_funds。
为什么必须是全量 snapshot, 而不是把 position_update 攒成一份视图
------------------------------------------------------------------
§5.8 `position_update` **增量**, 有变化即推拿增量攒视图有两个致命处:
漏一条推送就静默错账 (而账本错了看不出来); 更麻烦的是某只票清零这类事件, 对端
若不推 qty=0 就根本无法表达 视图里那只票会永远挂着对账要的是"此刻账户里到底
有什么"这个闭集, 只有全量快照给得出。增量只用来提示"有变动了", 不作为对账事实源。
握手后**立刻先拉一次**, 不等第一个间隔: 重建账本时没人愿意等 5 分钟, 而且刚连上
正是最需要知道对端持仓的时刻
发送失败不抛 (除连接类异常) 快照是兜底手段, 拉不到只是让 reconcile 退回表通道
并告警, 不该把整条通道拖down故障即守成是不产生新指令, 不是一有毛病就断线
"""
first = True
while not self._stop.is_set():
if not first:
with contextlib.suppress(asyncio.TimeoutError):
await asyncio.wait_for(self._stop.wait(),
timeout=self._p("query_interval_sec", 300))
if self._stop.is_set():
return
first = False
try:
await self._send(wsc.T_Q_POS, {}, seed)
await self._send(wsc.T_Q_FUNDS, {}, seed)
except asyncio.CancelledError:
raise
except Exception as e:
if _is_conn_error(e):
raise
logger.warning("快照查询发送失败 (下轮继续): %s", _brief_err(e))
async def _outbox_loop(self, seed: str): async def _outbox_loop(self, seed: str):
"""出口出栈: QUEUED → 签名 → send → SENT。顺带处理待撤单。 """出口出栈: QUEUED → 签名 → send → SENT。顺带处理待撤单。

View File

@ -101,6 +101,7 @@ class Settings(BaseSettings):
PMS_QMT_HEARTBEAT_DB_SEC: int = 2 # ws 进程写存活心跳的间隔 PMS_QMT_HEARTBEAT_DB_SEC: int = 2 # ws 进程写存活心跳的间隔
PMS_QMT_HEARTBEAT_STALE_SEC: int = 15 # 心跳陈旧超此秒数 → 判定进程已死, dispatcher 拒发 PMS_QMT_HEARTBEAT_STALE_SEC: int = 15 # 心跳陈旧超此秒数 → 判定进程已死, dispatcher 拒发
PMS_QMT_CONNECT_TIMEOUT_SEC: int = 10 # 建连超时 PMS_QMT_CONNECT_TIMEOUT_SEC: int = 10 # 建连超时
PMS_QMT_QUERY_INTERVAL_SEC: int = 300 # 全量快照轮询 query_positions+query_funds (协议 §6.2)
PMS_QMT_SEND_MAX_ATTEMPTS: int = 3 # 单张委托发送重试上限, 试满置 SEND_FAILED PMS_QMT_SEND_MAX_ATTEMPTS: int = 3 # 单张委托发送重试上限, 试满置 SEND_FAILED
# 下面两项是**密钥**: 只从 .env 注入, 不入库、不进 ParamStore、不上页面 (协议 §10.1.1)。 # 下面两项是**密钥**: 只从 .env 注入, 不入库、不进 ParamStore、不上页面 (协议 §10.1.1)。
# param_store.SECRET_KEYS 已把它们挡在可调参数与页面快照之外。 # param_store.SECRET_KEYS 已把它们挡在可调参数与页面快照之外。
@ -144,6 +145,10 @@ class Settings(BaseSettings):
# --- 对账与回放 --- # --- 对账与回放 ---
PMS_REPLAY_INTERVAL_MIN: int = 5 PMS_REPLAY_INTERVAL_MIN: int = 5
PMS_RECON_ALARM_DAYS: int = 3 # 连续不一致 N 日升级 ERROR PMS_RECON_ALARM_DAYS: int = 3 # 连续不一致 N 日升级 ERROR
# 持仓事实源: ws_first(默认, ws 快照为主/表为兜底) / ws_only / table_only
# 目标架构下 trading_service 全量退出业务, 那张表没有明确写入方; 协议 §6.2 定的也是 ws 快照。
PMS_RECON_SOURCE: str = "ws_first"
PMS_RECON_WS_SNAPSHOT_MAX_AGE_SEC: int = 900 # ws 快照超此秒数视为过期 (轮询 300s + 余量)
settings = Settings() settings = Settings()

View File

@ -367,6 +367,65 @@ def run():
assert wsc.T_ORDER_UPDATE not in wsc.NEEDS_LEDGER, \ assert wsc.T_ORDER_UPDATE not in wsc.NEEDS_LEDGER, \
"order_update 只作状态跟踪 —— 拿它入账会和逐笔 trade 重复计数" "order_update 只作状态跟踪 —— 拿它入账会和逐笔 trade 重复计数"
print("\n[F2] 持仓快照解析 (§5.7, 对账的事实源)")
@case("协议正文字段名: 归一成与 trading_position 同形")
def _():
r = wsc.parse_positions_snapshot({
"kind": "positions", "as_of": 1769500000000,
"items": [{"ts_code": "600000.SH", "total_qty": 6000, "avail_qty": 4000,
"frozen_qty": 0, "cost_price": 11.80, "market_price": 12.34}]})
eq(r["raw_count"], 1)
eq(r["as_of"], 1769500000000)
eq(r["rows"][0], {"ts_code": "600000.SH", "qty": 6000, "avail_qty": 4000,
"cost": 11.80, "frozen": 0, "price": 12.34})
eq(r["columns"]["qty"], "total_qty")
@case("券商原生字段名也认 (stock_code / total_quantity / positions)")
def _():
# 对端的实现多半是从券商结构体映射来的, 字段名写死就整批读不到 ——
# 而"读不到"与"真的没持仓"在对账那头长成同一个样子, 这是最不能混的两件事
r = wsc.parse_positions_snapshot({
"kind": "positions",
"positions": [{"stock_code": "sh600000", "total_quantity": "1100",
"available_quantity": 0, "cost_price": "9.273"}]})
eq(r["rows"][0]["ts_code"], "SH600000")
eq(r["rows"][0]["qty"], 1100)
eq(r["rows"][0]["cost"], 9.273)
eq(r["columns"]["qty"], "total_quantity")
@case("有条目却认不出数量列 → columns['qty'] 为 None (上层据此判「读不到」)")
def _():
r = wsc.parse_positions_snapshot({
"kind": "positions", "items": [{"ts_code": "600000.SH", "shengyu": 600}]})
eq(r["raw_count"], 1)
assert r["columns"]["qty"] is None, r["columns"]
eq(r["rows"][0]["qty"], None)
@case("空快照 / 缺 items / 脏条目都不崩")
def _():
eq(wsc.parse_positions_snapshot({"kind": "positions", "items": []})["rows"], [])
bad = wsc.parse_positions_snapshot({"kind": "positions", "as_of": 1})
eq(bad["rows"], [])
assert bad["columns"]["qty"] is None
assert "kind" in bad["unknown_fields"], bad["unknown_fields"]
eq(wsc.parse_positions_snapshot(None)["rows"], [])
# 非 dict 条目与空代码一律跳过, 不带崩整批
r = wsc.parse_positions_snapshot({"kind": "positions", "items": [
"junk", {"ts_code": "", "total_qty": 100},
{"ts_code": "600000.SH", "total_qty": 100}]})
eq(len(r["rows"]), 1)
@case("资金快照解析 (data 内层与平铺都认)")
def _():
r = wsc.parse_funds_snapshot({"kind": "funds", "as_of": 9, "data": {
"total_asset": 2013456.78, "available_cash": 312450.10,
"sell_return_today": 48900.00}})
eq(r["total_asset"], 2013456.78)
eq(r["sell_return_today"], 48900.00)
eq(r["frozen_cash"], None)
eq(wsc.parse_funds_snapshot({"total_asset": 1.0})["total_asset"], 1.0)
print("\n[I] ws 逐笔成交 → 账本动作 (清单 4/5)") print("\n[I] ws 逐笔成交 → 账本动作 (清单 4/5)")
from app.core import recon as rc from app.core import recon as rc

View File

@ -327,6 +327,7 @@ class FakeQmtRepo:
"resync_flag": 0, "last_error": None, "stat": {}} "resync_flag": 0, "last_error": None, "stat": {}}
self.orders = {} self.orders = {}
self.inbox = {} # seq → 上行消息行 (pms_qmt_inbox) self.inbox = {} # seq → 上行消息行 (pms_qmt_inbox)
self.snapshots = [] # snapshot 上行 (对账的事实源, §5.7)
def set_channel(self, conn_state="ONLINE", alive=True): def set_channel(self, conn_state="ONLINE", alive=True):
"""摆一个通道状态出来 —— 「进程活没活」和「连接通没通」是两件事。""" """摆一个通道状态出来 —— 「进程活没活」和「连接通没通」是两件事。"""
@ -399,6 +400,24 @@ class FakeQmtRepo:
def inbox_orphan_count(self): def inbox_orphan_count(self):
return sum(1 for r in self.inbox.values() if int(r["processed"]) == 3) return sum(1 for r in self.inbox.values() if int(r["processed"]) == 3)
def put_snapshot(self, *, kind="positions", items=None, data=None, age_sec=0, seq=None):
"""塞一条 snapshot 上行 (§5.7)。age_sec 用来造"快照过期""""
seq = int(seq if seq is not None else (max(self.inbox) + 1 if self.inbox else 1000))
pl = {"kind": kind, "as_of": 1769500000000}
if items is not None:
pl["items"] = items
if data is not None:
pl["data"] = data
self.snapshots.append({"seq": seq, "payload": pl, "age_sec": float(age_sec),
"received_at": datetime.now()})
return pl
def latest_snapshot(self, kind, *, scan=60):
for s in sorted(self.snapshots, key=lambda x: x["seq"], reverse=True):
if str(s["payload"].get("kind") or "").lower() == str(kind).lower():
return dict(s)
return None
def install_fakes(prices=None, positions=None, params=None, high5=None): def install_fakes(prices=None, positions=None, params=None, high5=None):
"""把内存桩装到各模块上, 返回 FakeRepo 实例 (ws 通道桩挂在 .qmt 上)。""" """把内存桩装到各模块上, 返回 FakeRepo 实例 (ws 通道桩挂在 .qmt 上)。"""
@ -409,7 +428,7 @@ def install_fakes(prices=None, positions=None, params=None, high5=None):
fake.qmt = FakeQmtRepo() fake.qmt = FakeQmtRepo()
for _n in ("get_state", "inbox_pending_count", "queue_depth", "enqueue_order", for _n in ("get_state", "inbox_pending_count", "queue_depth", "enqueue_order",
"request_cancel", "get_order", "inbox_pending", "inbox_mark", "request_cancel", "get_order", "inbox_pending", "inbox_mark",
"inbox_orphan_count"): "inbox_orphan_count", "latest_snapshot"):
setattr(qmt_repo, _n, getattr(fake.qmt, _n)) setattr(qmt_repo, _n, getattr(fake.qmt, _n))
fake.params.update(params or {}) fake.params.update(params or {})
for p in (positions or []): for p in (positions or []):
@ -1111,6 +1130,94 @@ def _():
assert fake.instructions["INS_A2"]["status"] == "EXPIRED" assert fake.instructions["INS_A2"]["status"] == "EXPIRED"
@case("对账事实源·ws 快照新鲜 → 以 ws 为准建账 (表空集不再当事实)")
def _():
from app.services import ledger_service as ls
fake = install_fakes(prices={"600000.SH": 10.0})
fake.qmt.put_snapshot(items=[{"ts_code": "600000.SH", "total_qty": 1100,
"avail_qty": 1100, "cost_price": 9.273}])
src = ls.positions_source()
assert src["source"] == ls.SRC_WS, src
assert src["rows"][0]["qty"] == 1100 and src["rows"][0]["cost"] == 9.273, src["rows"]
assert not src["alerts"], src["alerts"] # 表是空集, 不算"不一致"
# 账本空 + ws 有持仓 = 首次建账, 爆炸半径闸放行, 成本价取下游 cost_price
r = ls.reconcile(apply_fix=True)
assert r["source"] == "ws" and not r.get("blocked"), r
assert fake.positions["600000.SH"]["total_qty"] == 1100, fake.positions
assert abs(float(fake.positions["600000.SH"]["avg_cost"]) - 9.273) < 1e-6, \
"开仓价必须取下游 cost_price —— 拿现价当成本会让安全垫齐刷刷是 0"
@case("对账事实源·前缀式代码归一成点式 (不归一会造出一轮双向全量重写)")
def _():
from app.services import ledger_service as ls
fake = install_fakes(prices={"600000.SH": 10.0})
fake.qmt.put_snapshot(items=[{"stock_code": "SH600000", "total_quantity": 600,
"cost_price": 10.0}])
src = ls.positions_source()
assert src["rows"][0]["ts_code"] == "600000.SH", src["rows"]
@case("对账事实源·ws 与表对不上 → 用 ws 但必须留痕告警")
def _():
from app.repo import downstream_repo
from app.services import ledger_service as ls
fake = install_fakes(prices={"600000.SH": 10.0})
fake.qmt.put_snapshot(items=[{"ts_code": "600000.SH", "total_qty": 1100}])
downstream_repo.fetch_positions = lambda: {
"rows": [{"ts_code": "600000.SH", "qty": 700, "avail_qty": 700, "cost": 9.0,
"frozen": 0, "price": 10.0}],
"columns": {"qty": "total_quantity"}, "raw_count": 1}
src = ls.positions_source()
assert src["source"] == ls.SRC_WS, src["source"]
assert src["rows"][0]["qty"] == 1100, "以 ws 为准"
codes = [a["code"] for a in src["alerts"]]
assert "SOURCE_DISAGREE" in codes, src["alerts"]
@case("对账事实源·ws 快照过期 → 退回表并说明退回原因")
def _():
from app.repo import downstream_repo
from app.services import ledger_service as ls
fake = install_fakes(prices={"600000.SH": 10.0},
params={"PMS_RECON_WS_SNAPSHOT_MAX_AGE_SEC": "900"})
fake.qmt.put_snapshot(items=[{"ts_code": "600000.SH", "total_qty": 1100}], age_sec=5000)
downstream_repo.fetch_positions = lambda: {
"rows": [{"ts_code": "600000.SH", "qty": 700, "avail_qty": 700, "cost": 9.0,
"frozen": 0, "price": 10.0}],
"columns": {"qty": "total_quantity"}, "raw_count": 1}
src = ls.positions_source()
assert src["source"] == ls.SRC_TABLE, src["source"]
assert src["rows"][0]["qty"] == 700
msg = " ".join(a["message"] for a in src["alerts"])
assert "过期" in msg, src["alerts"]
@case("对账事实源·两个源都拿不到 + 本端有持仓 → **拒绝对账**, 一股都不许核销")
def _():
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() # 没塞快照, 表也空
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, "持仓不得被动过"
@case("对账事实源·两个源都空 + 账本也空 → 不报 ERROR (刚清账等对端装持仓)")
def _():
from app.services import ledger_service as ls
install_fakes()
r = ls.reconcile(apply_fix=True)
# 每分钟一跳的轻对账在这个状态下会反复走到这儿, 刷 ERROR 只会把真告警埋掉
assert not r.get("blocked") and r["severity"] == "OK", r
assert "无可对之账" in (r.get("note") or ""), r
@case("ws 入账·三分流: 真单入账 / 联调单跳过 / **孤儿成交挂起不入账**") @case("ws 入账·三分流: 真单入账 / 联调单跳过 / **孤儿成交挂起不入账**")
def _(): def _():
from app.repo import qmt_repo from app.repo import qmt_repo

View File

@ -70,6 +70,18 @@ def cmd_status(args):
print(f" 待入账上行 {ch['inbox_pending']}" print(f" 待入账上行 {ch['inbox_pending']}"
+ (f" · 挂起待人工 {ch['orphan_held']}" if ch.get("orphan_held") else "")) + (f" · 挂起待人工 {ch['orphan_held']}" if ch.get("orphan_held") else ""))
print(f" 出口队列 {ch['queue'] or ''}") print(f" 出口队列 {ch['queue'] or ''}")
# 快照是对账的事实源 (协议 §6.2), 拿不到就意味着账本建不起来 —— 必须一眼能看见
try:
snap = qmt_repo.latest_snapshot("positions")
except Exception as e:
snap = None
print(f" 持仓快照 读取失败: {type(e).__name__}: {e}")
if snap:
items = (snap["payload"] or {}).get("items")
print(f" 持仓快照 {len(items) if isinstance(items, list) else '?'} 个条目 · "
f"{snap['age_sec']:.0f}s 前 (seq={snap['seq']})")
elif snap is None:
print(f" 持仓快照 **从未收到** (对账没有事实源, 账本建不起来)")
s = st.get("stat") or {} s = st.get("stat") or {}
if s: if s:
print(f" 收发计数 收 {s.get('rx')}{s.get('tx')} · 成交 {s.get('trades')}" print(f" 收发计数 收 {s.get('rx')}{s.get('tx')} · 成交 {s.get('trades')}"