From e8902c21e8b34527611a247a94bace9953b5bc2d Mon Sep 17 00:00:00 2001 From: claude_dev Date: Wed, 7 Oct 2026 13:19:18 +0800 Subject: [PATCH] =?UTF-8?q?fix(monitor):=20account=20=E7=BE=A4=E9=95=BF?= =?UTF-8?q?=E5=B0=BE=E5=85=AD=E4=BF=AE+=E4=B8=A4=E4=BB=B6=E5=86=BB?= =?UTF-8?q?=E7=BB=93(=E7=BB=88=E8=A3=81=C2=A73.2,=E6=A0=87=E5=B0=BA=3D?= =?UTF-8?q?=E5=BA=94=E6=80=A5=E7=8F=AD=E8=BD=A6=E4=B8=8D=E6=B1=A1=E6=9F=93?= =?UTF-8?q?10-08=E5=BD=92=E5=9B=A0)=E2=80=94=E2=80=94=E5=85=AD=E4=BB=B6?= =?UTF-8?q?=E5=85=A8neutral:=20connect=E9=94=99=E8=AF=AF=E5=88=86=E6=94=AF?= =?UTF-8?q?=E5=8D=95=E6=AC=A1=E8=B0=83=E7=94=A8/=E6=97=A0mini=5Fpath?= =?UTF-8?q?=E5=91=8A=E8=AD=A6=E8=8A=82=E6=B5=81(=E9=A6=96=E7=8E=B0+30?= =?UTF-8?q?=E8=BD=AE=E5=BF=83=E8=B7=B3)/probe=E4=BC=9A=E8=AF=9Did=E6=B4=BE?= =?UTF-8?q?=E7=94=9Fper-process(base+pid=E9=98=B2=E9=87=8D=E5=90=AF?= =?UTF-8?q?=E6=92=9Esession)/=E5=90=8C=E6=AD=A5=E6=9F=A5=E8=AF=A2=E8=B6=85?= =?UTF-8?q?=E6=97=B6=E6=8A=A4=E6=A0=8F30s(daemon=E7=BA=BF=E7=A8=8Bjoin,?= =?UTF-8?q?=E8=B6=85=E6=97=B6=3D=E5=BD=93=E8=BD=AE=E5=A4=B1=E8=B4=A5+?= =?UTF-8?q?=E9=87=8D=E7=BD=AE,bs.login=E6=8C=82=E6=AD=BB=E5=90=8C=E6=97=8F?= =?UTF-8?q?,fail-closed=E8=AF=AD=E4=B9=89=E4=B8=8D=E5=8F=98)/StockAccount?= =?UTF-8?q?=E6=98=BE=E5=BC=8Faccount=5Ftype/=E5=BF=AB=E7=85=A7=E5=AD=97?= =?UTF-8?q?=E6=AE=B5=E7=9C=9F=E7=BC=BA=E5=A4=B1=E9=A6=96=E7=8E=B0WARN?= =?UTF-8?q?=E6=8C=890(=E9=9D=99=E9=BB=98=E9=9B=B6=E5=80=BC=E4=B8=8E?= =?UTF-8?q?=E7=9C=9F=E9=9B=B6=E5=8F=AF=E8=BE=A8,=E5=80=BC=E8=AF=AD?= =?UTF-8?q?=E4=B9=89=E4=B8=8D=E5=8A=A8);=20=E5=86=BB=E7=BB=93:=20traded=5F?= =?UTF-8?q?at=20LIKE=E5=AE=BD=E5=8C=B9=E9=85=8D=3Dfills=20v1.2=E9=A6=96?= =?UTF-8?q?=E5=AE=9E=E8=AF=84=E5=8F=96=E6=95=B0=E9=9D=A2=E7=AA=97=E5=8F=A3?= =?UTF-8?q?=E5=90=8E=E5=86=8D=E5=8A=A8/wmic=E6=97=B6=E5=8C=BA=3D=E7=A1=AE?= =?UTF-8?q?=E8=AE=A4=E5=85=B3=E5=8D=95(=E4=B8=89=E6=97=A5=E5=88=A4?= =?UTF-8?q?=E8=AF=BB=E7=BB=BF=E5=AE=9E=E8=AF=81);=20spec=20=C2=A79.A-19?= =?UTF-8?q?=E5=90=8Cpush;=20=E6=B5=8B=E8=AF=95=5Fpatch=5Fqmt=E8=A1=A5?= =?UTF-8?q?=E7=B1=BB=E5=B1=9E=E6=80=A7=E5=A4=8D=E4=BD=8D(=E8=A3=B8setattr?= =?UTF-8?q?=E8=B7=A8=E6=B5=8B=E6=B3=84=E6=BC=8F)=20[vps]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../specs/2026-10-03-monitoring-design.md | 2 + sanguo_live/account_monitor.py | 107 ++++++++++++--- tests/api/test_account_monitor.py | 126 +++++++++++++++++- 3 files changed, 210 insertions(+), 25 deletions(-) diff --git a/docs/superpowers/specs/2026-10-03-monitoring-design.md b/docs/superpowers/specs/2026-10-03-monitoring-design.md index 37830977..156c77ee 100644 --- a/docs/superpowers/specs/2026-10-03-monitoring-design.md +++ b/docs/superpowers/specs/2026-10-03-monitoring-design.md @@ -422,6 +422,8 @@ kind 语义(检查分发依据,check_entry data_monitor.py:56-232): **A-18 QMT 账号硬编码扩散(10-05 用户拍板 ⚖️-2 / P2-18 根治)**:现象=模拟柜台账号 66639661 曾硬编码 32 文件 88 处(灾备腿纯常量/桥腿 env 带默认/wrapper 注入字面量/env 名分裂三个),换号需动 14 生产文件+≥4 生产 wrapper。根因三层=schtask 无 per-task env+机器级 env 未设(wrapper 被迫字面量)/「部署身份」无单一权威源(三代各自发明)/模拟号零安全压力致 copy-paste 传播(run_shadow_bridge.ps1 同文件 redis 口令纪律在线、账号却是字面量=铁证)。修法=config/qmt_identity.json 单一权威源(JSON 非 yaml:灾备腿 stdlib 纯净不能依赖 pyyaml)+gate_common.load_identity 正典(:14-59)+各腿镜像+CI 同步测+码面字面量 grep-pin(tests/trader/test_qmt_identity.py,只 config 三件套可含字面量)。判别复用=换号只改 json+两 yaml(live.yaml/data_platform.yaml 由同步测钉一致)。全案=audit/20261005_qmt_identity_env_plan.md;生产 wrapper 字面量(库外 C:\sanguo_bigqmt)10-08 窗口后按非代码资产流程收尾。 +**A-19 account 群长尾六修+两件冻结(10-07 长尾批,终裁 §3.2;标尺=应急班车捎带部署不污染 10-08 窗口归因)**:account_monitor 系被动采集器(快照喂 console/API,不进任何灯色判定)六件全 neutral 当班修(sanguo_live/account_monitor.py,254→317 行):①connect 错误分支单次调用(旧二次调 `t.connect()` 拼消息=副作用+二次值可漂)②无 mini_path 跳过告警节流(首现+每 30 轮心跳;79c6b8b 只节流了配置读取失败)③probe 会话 id 派生 per-process(`_derive_session_id`=base+pid mod,进程内稳定跨进程不撞,防重启与残留旧连接撞 session)④同步查询超时护栏 `_call_with_timeout` 30s(daemon 线程 join;超时=当轮失败 note_failure+重置连接,bs.login 挂死同族防御,fail-closed 语义不变)⑤`StockAccount(account,"STOCK")` 显式账户类型(=现默认零行为差,防上游 API 漂移改默认)⑥快照字段真缺失(attr 缺失≠真零)`_num` 每进程首现 WARN 后按 0(静默零值与真零可辨,值语义不动)。**冻结两件**:`traded_at LIKE` 宽匹配=**10-08 21:30 fills v1.2 首实评取数面(strategy_checks.py:155 `_day_fills`),任何口径变动毁首评归因**→窗口绿后再修(届时有真班基线可对拍);wmic 时区假设=确认关单零代码(连续三日 wmic_shadow=3≡db_running=3 判读绿实证+§4.2-3 CreationDate 保守向留痕在位)。测试=tests/api/test_account_monitor.py 六件契约+`_patch_qmt` 补类属性复位(裸 setattr 跨测试泄漏防护)。 + ### §9.B 数据面(data,25 条;格式=现象/根因/判别法/修法坐标) 1. **mirror 假黄连三日(0926~28)**:push.log DONE 行只带 [HH:MM:SS] 不带日期(仅 start 行带)→旧 `want_day in line` 过滤永剩 start 行。判别=DONE 行格式先实证。修=start 段归属法(nas_monitor_check.py:166-185)。 diff --git a/sanguo_live/account_monitor.py b/sanguo_live/account_monitor.py index 5d4d3aaa..3864c490 100644 --- a/sanguo_live/account_monitor.py +++ b/sanguo_live/account_monitor.py @@ -23,13 +23,57 @@ import yaml logger = logging.getLogger(__name__) -# 专属 probe 会话 id:int,量级刻意远离 bullet_trade 默认的 int(time*1000)(~1.7e12), +# 专属 probe 会话 id 基数:int,量级刻意远离 bullet_trade 默认的 int(time*1000)(~1.7e12), # 不与引擎连接撞 session。 PROBE_SESSION_ID = 880811 _DEFAULT_INTERVAL_SEC = 60.0 _HEARTBEAT_EVERY_N_FAILURES = 30 +_QUERY_TIMEOUT_SEC = 30.0 # 同步查询超时护栏(长尾④,bs.login 挂死同族) +_MISSING = object() # attr 缺失哨兵(≠值为 None/0, 长尾⑥) +_MISSING_WARNED: set = set() # (类名, 字段) 每进程只首现告警一次 + + +def _derive_session_id() -> int: + """probe 会话 id 派生:进程内稳定、跨进程不撞——固定 id 在进程重启后与 + 残留旧连接撞 session(长尾③);base+pid 抖动,量级仍远离引擎 int(time*1000).""" + return PROBE_SESSION_ID + os.getpid() % 100000 + + +def _call_with_timeout(fn, *args): + """同步查询超时护栏(长尾④):超时=TimeoutError→当轮失败(fail-closed 语义 + 不变,「无限挂」变「超时重试」);守护线程泄漏有界,不阻塞后续轮次。""" + box: dict[str, Any] = {} + + def _run() -> None: + try: + box["ret"] = fn(*args) + except BaseException as e: # noqa: BLE001 + box["err"] = e + + t = threading.Thread(target=_run, daemon=True) + t.start() + t.join(_QUERY_TIMEOUT_SEC) + if t.is_alive(): + raise TimeoutError(f"query timeout {_QUERY_TIMEOUT_SEC}s") + if "err" in box: + raise box["err"] + return box.get("ret") + + +def _num(obj: Any, field: str) -> float: + """快照数值读取(长尾⑥):attr 真缺失(≠值为 None/0)→每进程首现 WARN 后按 0 + ——静默零值与真零可辨;值语义不动(None/缺省照旧折 0)。""" + v = getattr(obj, field, _MISSING) + if v is _MISSING: + key = (type(obj).__name__, field) + if key not in _MISSING_WARNED: + _MISSING_WARNED.add(key) + logger.warning("[account-monitor] 快照字段缺失 %s.%s 按 0 记", *key) + return 0.0 + return float(v or 0) + _CFG_FAILURES = 0 # watch 配置读取失败计数(节流告警用, _note_failure 同哲学) @@ -87,18 +131,21 @@ class AccountMonitor(threading.Thread): self, db_path: str, interval_sec: float = _DEFAULT_INTERVAL_SEC, - session_id: int = PROBE_SESSION_ID, + session_id: int | None = None, extra_accounts: dict[str, str] | None = None, ) -> None: super().__init__(daemon=True, name="account-monitor") self.db_path = db_path self.interval_sec = interval_sec - self.session_id = session_id + # None=进程内派生(跨进程不撞,长尾③);显式传值=测试/运维覆盖 + self.session_id = session_id if session_id is not None \ + else _derive_session_id() # None=首 poll 时读 config;{}=禁用(测试用) self.extra_accounts = extra_accounts self._stop_event = threading.Event() self._traders: dict[str, Any] = {} # mini_path → XtQuantTrader self._fail_counts: dict[str, int] = {} # mini_path → 连续失败计数 + self._skip_counts: dict[str, int] = {} # account → 无 mini_path 跳过计数(长尾②) # ---------------- 生命周期 ---------------- @@ -166,49 +213,64 @@ class AccountMonitor(threading.Thread): return written def _poll_account(self, account: str, mini_path: str, upsert: Any) -> bool: - """单账户采集。返回是否写入。mini_path 为空 → 告警跳过(连不上 QMT)。""" + """单账户采集。返回是否写入。mini_path 为空 → 节流告警跳过(连不上 QMT)。""" if not mini_path: - logger.warning( - "[account-monitor] 账户 %s 无 mini_path,跳过(需 live_accounts " - "行携带或快照 sticky 记录)", account) + n = self._skip_counts.get(account, 0) + 1 # 长尾②:节流 + self._skip_counts[account] = n + if n == 1 or n % _HEARTBEAT_EVERY_N_FAILURES == 0: + logger.warning( + "[account-monitor] 账户 %s 无 mini_path,跳过(需 live_accounts " + "行携带或快照 sticky 记录; 第%d次)", account, n) return False trader = self._ensure_trader(mini_path) if trader is None: return False XtQuantTrader, StockAccount = _import_qmt() - acc_obj = StockAccount(account) - asset = trader.query_stock_asset(acc_obj) + acc_obj = StockAccount(account, "STOCK") # 长尾⑤:显式账户类型 + try: # 长尾④:超时护栏 + asset = _call_with_timeout(trader.query_stock_asset, acc_obj) + except TimeoutError as e: + self._note_failure(mini_path, repr(e)) + self._reset_trader(mini_path) + return False if asset is None: self._note_failure(mini_path, "query_stock_asset None") self._reset_trader(mini_path) # 可能断连,下轮重建 return False - positions = trader.query_stock_positions(acc_obj) or [] + try: + positions = _call_with_timeout( + trader.query_stock_positions, acc_obj) or [] + except TimeoutError as e: + self._note_failure(mini_path, repr(e)) + self._reset_trader(mini_path) + return False rows = [] for p in positions: - vol = float(getattr(p, "volume", 0) or 0) + vol = _num(p, "volume") if vol <= 0: continue rows.append({ "symbol": str(getattr(p, "stock_code", "") or ""), "volume": vol, - "can_use": float(getattr(p, "can_use_volume", 0) or 0), - "avg_price": float(getattr(p, "avg_price", 0) or 0), - "mv": float(getattr(p, "market_value", 0) or 0), + "can_use": _num(p, "can_use_volume"), + "avg_price": _num(p, "avg_price"), + "mv": _num(p, "market_value"), }) + cash = _num(asset, "cash") + mv = _num(asset, "market_value") + total = _num(asset, "total_asset") upsert( self.db_path, account, - cash=float(getattr(asset, "cash", 0) or 0), - market_value=float(getattr(asset, "market_value", 0) or 0), - total=float(getattr(asset, "total_asset", 0) or 0), + cash=cash, + market_value=mv, + total=total, positions=rows, mini_path=mini_path, ) self._fail_counts.pop(mini_path, None) logger.info( "[account-monitor] 快照 %s: cash=%.0f mv=%.0f total=%.0f 持仓%d只", - account, float(getattr(asset, "cash", 0) or 0), - float(getattr(asset, "market_value", 0) or 0), - float(getattr(asset, "total_asset", 0) or 0), len(rows)) + account, cash, mv, total, len(rows)) return True # ---------------- 连接管理 ---------------- @@ -222,8 +284,9 @@ class AccountMonitor(threading.Thread): try: t = XtQuantTrader(mini_path, self.session_id) t.start() - if t.connect() not in (0, None): - raise RuntimeError(f"connect 返回 {t.connect()}") + rc = t.connect() # 长尾①:单次调用 + if rc not in (0, None): + raise RuntimeError(f"connect 返回 {rc}") self._traders[mini_path] = t return t except Exception as e: # noqa: BLE001 diff --git a/tests/api/test_account_monitor.py b/tests/api/test_account_monitor.py index a84b5f27..57d436ea 100644 --- a/tests/api/test_account_monitor.py +++ b/tests/api/test_account_monitor.py @@ -75,6 +75,7 @@ class _FakeTrader: """假 XtQuantTrader:类属性记录实例,可编程返回资产/持仓。""" instances = [] connect_ret = 0 + connect_calls = 0 asset = None positions = [] @@ -88,6 +89,7 @@ class _FakeTrader: return 0 def connect(self): + type(self).connect_calls += 1 return type(self).connect_ret def stop(self): @@ -104,12 +106,19 @@ class _FakeTrader: def _patch_qmt(monkeypatch, **attrs): + # 类属性默认先复位(裸 setattr 不随 monkeypatch 回滚,跨测试泄漏防护) + _FakeTrader.instances = [] + _FakeTrader.connect_calls = 0 + _FakeTrader.connect_ret = 0 + _FakeTrader.asset = None + _FakeTrader.positions = [] for k, v in attrs.items(): setattr(_FakeTrader, k, v) - _FakeTrader.instances = [] monkeypatch.setattr( am, "_import_qmt", - lambda: (_FakeTrader, lambda acc: SimpleNamespace(account=acc))) + lambda: (_FakeTrader, + lambda acc, account_type="STOCK": SimpleNamespace( + account=acc))) def test_watch_targets_union(db, monkeypatch): @@ -166,7 +175,7 @@ def test_poll_once_upserts_snapshot(db, monkeypatch): # 连接复用:同路径第二次 poll 不新建 trader assert m.poll_once() == 1 assert len(_FakeTrader.instances) == 1 - assert _FakeTrader.instances[0].session_id == am.PROBE_SESSION_ID + assert _FakeTrader.instances[0].session_id == am._derive_session_id() def test_poll_query_none_resets_and_no_write(db, monkeypatch, caplog): @@ -270,3 +279,114 @@ def test_extra_from_config_fail_safe_and_throttled(tmp_path, monkeypatch, caplog assert am._extra_from_config( cfg_path=str(tmp_path / "none.yaml")) == {} # 文件缺失 assert len(caplog.records) == 1 # 两次失败只首败告警(节流契约) + + +# ---------------- 10-07 长尾批(monitoring 终裁 §3.2 account 群六件) ---------------- + +def test_connect_single_call_in_error_path(db, monkeypatch): + """长尾①: connect 失败分支不得重调 connect() 拼错误消息(二次调用=副作用+值可漂).""" + _patch_qmt(monkeypatch, connect_ret=-1, asset=None, positions=[]) + lp.upsert_account_snapshot(db, "A1", 1, 0, 1, [], "P1") + m = am.AccountMonitor(db, extra_accounts={}) + assert m.poll_once() == 0 + assert _FakeTrader.connect_calls == 1 + + +def test_empty_mini_path_warn_throttled(db, monkeypatch, caplog): + """长尾②: 无 mini_path 跳过告警须节流——首现 WARN,后续轮静默 + (79c6b8b 只节流了配置读取失败,此枝每轮刷屏).""" + _patch_qmt(monkeypatch, asset=None, positions=[]) + with sqlite3.connect(db) as conn: + conn.execute( + "INSERT INTO qmt_account_snapshot (account, mini_path, cash," + " market_value, total, positions, updated_at)" + " VALUES ('NOPATH','',1,0,1,'[]','x')") + conn.commit() + m = am.AccountMonitor(db, extra_accounts={}) + with caplog.at_level("WARNING"): + for _ in range(3): + assert m.poll_once() == 0 + assert len([r for r in caplog.records + if "无 mini_path" in r.message]) == 1 + + +def test_session_id_derived_per_process(db, monkeypatch): + """长尾③: probe 会话 id 派生=进程内稳定、跨进程不撞(固定 id 在进程重启后 + 与残留旧连接撞 session);量级仍远离引擎 int(time*1000).""" + _patch_qmt(monkeypatch, + asset=SimpleNamespace(cash=1.0, market_value=0.0, + total_asset=1.0)) + lp.upsert_account_snapshot(db, "A1", 1, 0, 1, [], "P1") + m = am.AccountMonitor(db, extra_accounts={}) + assert m.poll_once() == 1 + sid = _FakeTrader.instances[0].session_id + assert sid == am._derive_session_id() == am._derive_session_id() + assert isinstance(sid, int) and sid >= am.PROBE_SESSION_ID + + +def test_query_timeout_fail_closed(db, monkeypatch, caplog): + """长尾④: 同步查询挂死→超时护栏=当轮失败+重置连接,不挂 poll 线程 + (bs.login 挂死同族防御;fail-closed 语义不变).""" + import time + + class _Hanging(_FakeTrader): + def query_stock_asset(self, acc): + time.sleep(5) + return None + + _patch_qmt(monkeypatch, asset=None, positions=[]) + monkeypatch.setattr( + am, "_import_qmt", + lambda: (_Hanging, + lambda acc, account_type="STOCK": SimpleNamespace( + account=acc))) + monkeypatch.setattr(am, "_QUERY_TIMEOUT_SEC", 0.05) + lp.upsert_account_snapshot(db, "A1", 1, 0, 1, [], "P1") + m = am.AccountMonitor(db, extra_accounts={}) + t0 = time.monotonic() + with caplog.at_level("WARNING"): + assert m.poll_once() == 0 + assert time.monotonic() - t0 < 3 # 未被 5s sleep 拖住 + assert _FakeTrader.instances[-1].stopped # 超时即重置连接 + assert any("timeout" in r.message.lower() for r in caplog.records) + + +def test_stock_account_type_explicit(db, monkeypatch): + """长尾⑤: StockAccount 显式携带 account_type('STOCK'=现默认,零行为差) + ——防上游 API 漂移改默认时探针静默换了账户类型.""" + made = [] + + def _mk(acc, account_type="UNSET"): + made.append((acc, account_type)) + return SimpleNamespace(account=acc) + + _patch_qmt(monkeypatch, + asset=SimpleNamespace(cash=1.0, market_value=0.0, + total_asset=1.0)) + monkeypatch.setattr(am, "_import_qmt", lambda: (_FakeTrader, _mk)) + lp.upsert_account_snapshot(db, "A1", 1, 0, 1, [], "P1") + m = am.AccountMonitor(db, extra_accounts={}) + assert m.poll_once() == 1 + assert made and made[0][1] == "STOCK" + + +def test_missing_attr_visible_not_silent_zero(db, monkeypatch, caplog): + """长尾⑥: 持仓对象真缺字段(attr 缺失≠真零)→快照值仍按 0(语义不动)但 + 每进程首现 WARN——静默零值与真零可辨.""" + _patch_qmt(monkeypatch, + asset=SimpleNamespace(cash=1.0, market_value=0.0, + total_asset=1.0), + positions=[SimpleNamespace(stock_code="600036.SH", + volume=100)]) # 缺 can_use/avg/mv + lp.upsert_account_snapshot(db, "A1", 1, 0, 1, [], "P1") + m = am.AccountMonitor(db, extra_accounts={}) + with caplog.at_level("WARNING"): + assert m.poll_once() == 1 + n1 = len([r for r in caplog.records if "字段缺失" in r.message]) + assert n1 == 3 + with caplog.at_level("WARNING"): + assert m.poll_once() == 1 + assert len([r for r in caplog.records if "字段缺失" in r.message]) == n1 + snap = lp.get_account_snapshot(db, "A1") + assert snap["positions"][0]["volume"] == 100.0 + assert snap["positions"][0]["avg_price"] == 0.0 # 值语义不动