diff --git a/app/core/ws_codec.py b/app/core/ws_codec.py index bb13ed9..9f0d250 100644 --- a/app/core/ws_codec.py +++ b/app/core/ws_codec.py @@ -425,6 +425,106 @@ def dedup_key(type_: str, payload: dict): 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: """§5.5: amount 应等于 price × qty, 差异 > 0.01 元告警。""" try: diff --git a/app/repo/qmt_repo.py b/app/repo/qmt_repo.py index d43286c..7a9ab8a 100644 --- a/app/repo/qmt_repo.py +++ b/app/repo/qmt_repo.py @@ -281,6 +281,31 @@ def inbox_orphan_count() -> int: 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: """seq 之上的上行消息按 类型×入账状态 汇总。给 rewind 判"删了会不会出事"用。""" return fetch_all( diff --git a/app/services/ledger_service.py b/app/services/ledger_service.py index 6287bad..baa39aa 100644 --- a/app/services/ledger_service.py +++ b/app/services/ledger_service.py @@ -22,6 +22,7 @@ from app.core import command_spec as cs from app.core import cushion as cu from app.core import recon as rc 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.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 +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: - """账本 vs 下游持仓, 以下游为准修正并留痕。 + """账本 vs 下游持仓, 以下游为准修正并留痕。事实源见 positions_source()。 force=True 绕过爆炸半径限制 (见 _recon_blast_guard) —— 只在人工确认下游读数确实 正确之后使用。 """ out = {"ok": True, "diffs": [], "fixes": [], "errors": [], "columns": {}, "severity": rc.SEV_OK} - try: - ds = downstream_repo.fetch_positions() - except Exception as e: - out.update({"ok": False, "errors": [f"读 trading_position 失败: {type(e).__name__}: {e}"]}) - return out - out["columns"] = ds["columns"] - if ds["rows"] and ds["columns"].get("qty") is None: - out.update({"ok": False, "errors": [ - "下游持仓表未识别出数量列 —— 请按 QMT_INTERFACE_REQUIREMENTS A1/D1 取得 DDL 后, " - "把列名补进 downstream_repo.QTY_CANDIDATES"]}) + src = positions_source() + ds = {"rows": src["rows"], "columns": src["columns"], "raw_count": src["raw_count"]} + out.update({"source": src["source"], "source_mode": src["mode"], + "source_age_sec": src.get("age_sec"), "source_alerts": src["alerts"], + "columns": src["columns"]}) + for a in src["alerts"]: + (logger.error if a["level"] == "ERROR" else logger.warning)( + "[对账·事实源] %s", a["message"]) + + 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 # 下游行先过一遍代码合法性。券商表里混进非个股代码的原因很多 (联调测试单、B 股、 diff --git a/app/ws/runner.py b/app/ws/runner.py index 6d61ca2..e3ca092 100644 --- a/app/ws/runner.py +++ b/app/ws/runner.py @@ -106,6 +106,7 @@ class WsRunner: "beat_sec": param_store.get_int("PMS_QMT_HEARTBEAT_DB_SEC", 2), "max_attempts": param_store.get_int("PMS_QMT_SEND_MAX_ATTEMPTS", 3), "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: self._params = await _db(_load) @@ -115,7 +116,8 @@ class WsRunner: "url": settings.PMS_QMT_WS_URL, "heartbeat_sec": 5, "idle_timeout_sec": 15, "ack_batch": 20, "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)) @staticmethod @@ -287,7 +289,8 @@ class WsRunner: tasks = [asyncio.create_task(self._reader_loop(peer), name="reader"), asyncio.create_task(self._pinger_loop(seed), name="pinger"), 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 上, 光等它自己醒来会让 # 优雅退出白白拖满一个 idle 周期。谁先结束就收工, 剩下的直接 cancel。 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]) logger.info("[trade] %s %s %s股 @%s (待 worker 入账)", iid, 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: if not pl.get("ok"): # §5.6: 落库失败不影响成交事实, 但要留一条下游不一致告警 @@ -677,6 +695,41 @@ class WsRunner: return 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): """出口出栈: QUEUED → 签名 → send → SENT。顺带处理待撤单。 diff --git a/config/settings.py b/config/settings.py index 77beb2f..768b0b0 100644 --- a/config/settings.py +++ b/config/settings.py @@ -101,6 +101,7 @@ class Settings(BaseSettings): PMS_QMT_HEARTBEAT_DB_SEC: int = 2 # ws 进程写存活心跳的间隔 PMS_QMT_HEARTBEAT_STALE_SEC: int = 15 # 心跳陈旧超此秒数 → 判定进程已死, dispatcher 拒发 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 # 下面两项是**密钥**: 只从 .env 注入, 不入库、不进 ParamStore、不上页面 (协议 §10.1.1)。 # param_store.SECRET_KEYS 已把它们挡在可调参数与页面快照之外。 @@ -144,6 +145,10 @@ class Settings(BaseSettings): # --- 对账与回放 --- PMS_REPLAY_INTERVAL_MIN: int = 5 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() diff --git a/scripts/test_batch6_units.py b/scripts/test_batch6_units.py index bd227be..4462cbf 100644 --- a/scripts/test_batch6_units.py +++ b/scripts/test_batch6_units.py @@ -367,6 +367,65 @@ def run(): assert wsc.T_ORDER_UPDATE not in wsc.NEEDS_LEDGER, \ "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)") from app.core import recon as rc diff --git a/scripts/test_wiring.py b/scripts/test_wiring.py index 9f3ce53..655e6be 100644 --- a/scripts/test_wiring.py +++ b/scripts/test_wiring.py @@ -327,6 +327,7 @@ class FakeQmtRepo: "resync_flag": 0, "last_error": None, "stat": {}} self.orders = {} self.inbox = {} # seq → 上行消息行 (pms_qmt_inbox) + self.snapshots = [] # snapshot 上行 (对账的事实源, §5.7) def set_channel(self, conn_state="ONLINE", alive=True): """摆一个通道状态出来 —— 「进程活没活」和「连接通没通」是两件事。""" @@ -399,6 +400,24 @@ class FakeQmtRepo: def inbox_orphan_count(self): 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): """把内存桩装到各模块上, 返回 FakeRepo 实例 (ws 通道桩挂在 .qmt 上)。""" @@ -409,7 +428,7 @@ def install_fakes(prices=None, positions=None, params=None, high5=None): fake.qmt = FakeQmtRepo() for _n in ("get_state", "inbox_pending_count", "queue_depth", "enqueue_order", "request_cancel", "get_order", "inbox_pending", "inbox_mark", - "inbox_orphan_count"): + "inbox_orphan_count", "latest_snapshot"): setattr(qmt_repo, _n, getattr(fake.qmt, _n)) fake.params.update(params or {}) for p in (positions or []): @@ -1111,6 +1130,94 @@ def _(): 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 入账·三分流: 真单入账 / 联调单跳过 / **孤儿成交挂起不入账**") def _(): from app.repo import qmt_repo diff --git a/scripts/ws_smoke.py b/scripts/ws_smoke.py index 4fb364b..803691c 100644 --- a/scripts/ws_smoke.py +++ b/scripts/ws_smoke.py @@ -70,6 +70,18 @@ def cmd_status(args): print(f" 待入账上行 {ch['inbox_pending']} 条" + (f" · 挂起待人工 {ch['orphan_held']} 笔" if ch.get("orphan_held") else "")) 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 {} if s: print(f" 收发计数 收 {s.get('rx')} 发 {s.get('tx')} · 成交 {s.get('trades')}"