fix(live): EOD对账重启兜底归因——remark指纹重建本实例订单集+实例唯一label注入;09-01 002646终局实锤:引擎订单注册表=进程内存态(轮换重启即失,runtime持久化只恢复策略/账本),新进程18:49 EOD报「本实例=0」→无守恒缺口→记日,当日15:05后才补全的迟到成交永久失联(get_trades内部本就是query_stock_trades,重试环与直查同源都不缺接口,缺的是归因锚);修=①runner_live经live_config注入strategy_name=live_{id},订单remark(bt:live_17:hash)实例唯一(此前六实例适配文件同名,label全为live_strateg不可分;仅对注入后新单生效)②eod_reconcile补_own_orders_remarked:broker.get_orders(QMT当日订单,带remark)按前缀筛本实例→shim补进own_by_broker(setdefault,engine自己的单优先),状态复用broker._map_order_status,is_buy缺失弃用(宁漏勿错);_PENDING不重建(EOD兜底已覆盖当日全量,省一条风险面);+4重启形态测试+1注入缝钉子,557绿 [vps]
This commit is contained in:
@@ -24,6 +24,7 @@ from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from datetime import datetime
|
||||
from types import SimpleNamespace
|
||||
from typing import Any, Dict, Optional
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -206,6 +207,62 @@ def reconcile_pending(engine: Any, ledger: Any, db: str, account_id: int,
|
||||
return total
|
||||
|
||||
|
||||
# ------------------ 重启兜底:remark 指纹重建本实例订单集 ------------------
|
||||
def _own_orders_remarked(engine: Any, account_id: int) -> Dict[str, Any]:
|
||||
"""(2026-09-01 002646 事故)按 remark 指纹从 QMT 当日订单行重建本实例订单集。
|
||||
|
||||
engine 订单注册表是**进程内存态**——轮换/崩溃重启后 ``engine.get_orders()``
|
||||
为空,EOD 报「本实例=0」→ 无守恒缺口 → 记日,当日孤儿成交永久失去自动
|
||||
补插(当日 15:05 快照后视图才补全的迟到成交正是这形态)。QMT 当日订单行
|
||||
(``broker.get_orders`` → query_stock_orders)带 remark=``bt:{label}:{hash}``,
|
||||
label=runner_live 注入的 ``live_{account_id}``,跨重启可查——以它为兜底
|
||||
归因源,trades 行自身不带 remark,靠 order_id join 本表。
|
||||
|
||||
仅对注入后提交的订单生效(注入前的历史单 label 六实例共享不可分);
|
||||
``is_buy`` 缺失的行弃用——方向不明宁可漏归因,不可错归因(错方向入账
|
||||
会反向污染台账)。状态映射复用 broker 自己的 ``_map_order_status``,
|
||||
缺失时保 raw(EOD 守恒只认终态字符串,raw 码只致少计不误计)。
|
||||
"""
|
||||
broker = getattr(engine, "broker", None)
|
||||
getter = getattr(broker, "get_orders", None)
|
||||
if not callable(getter):
|
||||
return {}
|
||||
try:
|
||||
rows = getter() or []
|
||||
except Exception as e: # noqa: BLE001 - 兜底源失败不阻断主 EOD 流程
|
||||
logger.warning("[live-reconcile] remark 重建查 QMT 订单失败: %s", e)
|
||||
return {}
|
||||
prefix = f"bt:live_{account_id}:"
|
||||
mapper = getattr(broker, "_map_order_status", None)
|
||||
out: Dict[str, Any] = {}
|
||||
for row in rows:
|
||||
if not isinstance(row, dict):
|
||||
continue
|
||||
if not str(row.get("order_remark") or "").startswith(prefix):
|
||||
continue
|
||||
oid = str(row.get("order_id") or "")
|
||||
is_buy = row.get("is_buy")
|
||||
security = row.get("security")
|
||||
if not oid or is_buy is None or not security:
|
||||
logger.warning(
|
||||
"[live-reconcile] remark 重建弃用行(方向/代码缺失) order=%s", oid)
|
||||
continue
|
||||
status = row.get("status")
|
||||
if callable(mapper):
|
||||
try:
|
||||
mapped = mapper(status)
|
||||
status = getattr(mapped, "value", mapped)
|
||||
except Exception: # noqa: BLE001 - 映射失败保 raw
|
||||
pass
|
||||
out[oid] = SimpleNamespace(
|
||||
order_id=oid, _broker_order_id=oid, is_buy=bool(is_buy),
|
||||
amount=int(row.get("amount") or 0),
|
||||
filled=int(row.get("filled") or 0),
|
||||
status=status, security=str(security), add_time=None,
|
||||
)
|
||||
return out
|
||||
|
||||
|
||||
# ------------------ b 宽修:EOD 对账回填 ------------------
|
||||
def eod_reconcile(engine: Any, ledger: Any, db: str, account_id: int,
|
||||
strategy_name: str) -> Dict[str, int]:
|
||||
@@ -230,6 +287,10 @@ def eod_reconcile(engine: Any, ledger: Any, db: str, account_id: int,
|
||||
boid = getattr(o, "_broker_order_id", None)
|
||||
if boid:
|
||||
own_by_broker[str(boid)] = o
|
||||
# 重启兜底:engine 订单表是进程态,缺的按 remark 指纹从 QMT 当日订单补
|
||||
# (engine 自己的单信息更全,优先;setdefault 只填洞)
|
||||
for oid, shim in _own_orders_remarked(engine, account_id).items():
|
||||
own_by_broker.setdefault(oid, shim)
|
||||
broker = getattr(engine, "broker", None)
|
||||
trades_all = []
|
||||
getter = getattr(broker, "get_trades", None)
|
||||
|
||||
@@ -294,9 +294,17 @@ def run_live(provider_config: Dict[str, Any] | None = None) -> None:
|
||||
_instance_adapter(cfg["account_id"]),
|
||||
broker_factory=lambda: broker,
|
||||
# bullet_trade 每实例锁 runtime 目录(单实例设计);多实盘并行须各用独立目录
|
||||
live_config={"runtime_dir": str(
|
||||
Path(__file__).resolve().parent.parent / "runtime"
|
||||
/ f"live_{cfg['account_id'] or 'solo'}")},
|
||||
# strategy_name=实例唯一 remark label(2026-09-01 002646 事故:六实例适配
|
||||
# 文件同名 live_strategy.py → 订单 remark label 全为 "live_strateg" 跨实
|
||||
# 例不可分;EOD 对账重启后按 remark=bt:live_{id}: 重建本实例订单集,
|
||||
# 引擎订单表是进程态,只有 remark 跨重启可归因)。仅对注入后新单生效。
|
||||
live_config={
|
||||
"runtime_dir": str(
|
||||
Path(__file__).resolve().parent.parent / "runtime"
|
||||
/ f"live_{cfg['account_id'] or 'solo'}"),
|
||||
**({"strategy_name": f"live_{cfg['account_id']}"}
|
||||
if cfg["account_id"] else {}),
|
||||
},
|
||||
)
|
||||
logger.info(
|
||||
"组合 live engine 启动: strategy=%s max_pool=%s benchmark=%s cash=%s",
|
||||
|
||||
@@ -558,3 +558,89 @@ class TestLedgerSeenIds:
|
||||
assert snap == {"t1"}
|
||||
snap.add("t2") # 副本可改,不污染账本
|
||||
assert led.seen_trade_ids() == {"t1"}
|
||||
|
||||
|
||||
# ------------------ 重启兜底:remark 指纹重建(2026-09-01 002646 事故) ------------------
|
||||
def _qmt_order_row(order_id="1090", remark="bt:live_19:ab12cd34",
|
||||
security="600276.XSHG", amount=600, filled=600,
|
||||
status=56, is_buy=True):
|
||||
"""贴 QmtBroker.sync_orders 行键(状态为 raw 码,56=全部成交)。"""
|
||||
return {"order_id": order_id, "security": security, "amount": amount,
|
||||
"filled": filled, "status": status, "is_buy": is_buy,
|
||||
"order_remark": remark}
|
||||
|
||||
|
||||
def _restarted_engine(qmt_orders, qmt_trades):
|
||||
"""重启后的 engine 替身:订单注册表为空(进程态已失),broker 行为贴
|
||||
QmtBroker(get_orders 带 remark + _map_order_status 归一)。"""
|
||||
broker = SimpleNamespace(
|
||||
get_orders=lambda: list(qmt_orders),
|
||||
get_trades=lambda: list(qmt_trades),
|
||||
_map_order_status=lambda raw: SimpleNamespace(
|
||||
value="filled" if raw == 56 else f"raw_{raw}"),
|
||||
)
|
||||
return SimpleNamespace(get_orders=lambda: {}, broker=broker)
|
||||
|
||||
|
||||
class TestRemarkRebuildAfterRestart:
|
||||
"""002646 形态:轮换杀进程 → QMT 视图 18:48-19:15 才补全迟到成交 →
|
||||
新进程 EOD「本实例=0」记日,孤儿成交永久失联。remark 指纹重建 own 集
|
||||
后,重启照样能归因/补插/守恒。"""
|
||||
|
||||
def test_backfill_via_remark_after_restart(self, db):
|
||||
"""engine 订单空 + QMT 订单 remark=本实例 → 迟到成交照常补插+守恒平。"""
|
||||
led = LiveInstanceLedger(initial_cash=1_000_000)
|
||||
eng = _restarted_engine(
|
||||
[_qmt_order_row(order_id="1090", amount=600, filled=600)],
|
||||
[_qmt_trade(order_id="1090", security="600276.XSHG",
|
||||
amount=600, price=7.0, trade_id="99016")])
|
||||
summary = eod_reconcile(eng, led, db, 19, "s")
|
||||
assert summary["ours"] == 1 # remark 重建归因成功
|
||||
assert summary["backfilled"] == 1 # 002646 那笔 600 股补回来
|
||||
assert summary["conservation_gaps"] == []
|
||||
rows = live_reconcile_eod_rows(db, 19)
|
||||
assert rows and rows[-1]["volume"] == 600
|
||||
|
||||
def test_other_instance_remark_not_attributed(self, db):
|
||||
"""别家实例的 remark 前缀(live_18)不归因给 live_19。"""
|
||||
led = LiveInstanceLedger(initial_cash=1_000_000)
|
||||
eng = _restarted_engine(
|
||||
[_qmt_order_row(order_id="1080", remark="bt:live_18:ffff1111")],
|
||||
[_qmt_trade(order_id="1080", security="600276.XSHG",
|
||||
amount=600, price=7.0, trade_id="99017")])
|
||||
summary = eod_reconcile(eng, led, db, 19, "s")
|
||||
assert summary["ours"] == 0
|
||||
assert summary["backfilled"] == 0
|
||||
assert summary["foreign"] == 1
|
||||
|
||||
def test_row_without_side_excluded(self, db, caplog):
|
||||
"""is_buy 缺失的 remark 行弃用——方向不明宁可漏归因不可错归因。"""
|
||||
led = LiveInstanceLedger(initial_cash=1_000_000)
|
||||
row = _qmt_order_row(order_id="1091")
|
||||
row["is_buy"] = None
|
||||
eng = _restarted_engine(
|
||||
[row],
|
||||
[_qmt_trade(order_id="1091", security="600276.XSHG",
|
||||
amount=100, price=7.0, trade_id="99018")])
|
||||
with caplog.at_level("WARNING"):
|
||||
summary = eod_reconcile(eng, led, db, 19, "s")
|
||||
assert summary["backfilled"] == 0
|
||||
assert any("方向" in r.message or "弃用" in r.message
|
||||
for r in caplog.records)
|
||||
|
||||
def test_sell_direction_propagates(self, db):
|
||||
"""卖出方向的 remark 重建单,补插行 direction=sell(错方向入账=反向污染)。"""
|
||||
led = LiveInstanceLedger(initial_cash=1_000_000)
|
||||
eng = _restarted_engine(
|
||||
[_qmt_order_row(order_id="1092", is_buy=False)],
|
||||
[_qmt_trade(order_id="1092", security="600276.XSHG",
|
||||
amount=600, price=7.0, trade_id="99019")])
|
||||
summary = eod_reconcile(eng, led, db, 19, "s")
|
||||
assert summary["backfilled"] == 1
|
||||
rows = live_reconcile_eod_rows(db, 19)
|
||||
assert rows[-1]["direction"] == "sell"
|
||||
|
||||
|
||||
def live_reconcile_eod_rows(db, account_id):
|
||||
from sanguo_live.persistence import list_trades
|
||||
return list_trades(db, account_id)
|
||||
|
||||
@@ -56,3 +56,13 @@ def test_backoff_sleep_invoked_per_failure():
|
||||
with_connect_retry(b, attempts=5, wait_sec=7.0, _sleep=slept.append)
|
||||
b.connect()
|
||||
assert slept == [7.0, 7.0] # 每次失败后睡一次
|
||||
|
||||
|
||||
def test_live_config_carries_instance_strategy_name():
|
||||
"""(2026-09-01 002646)runner_live 经 live_config 注入 strategy_name=
|
||||
live_{id},使 bullet_trade 订单 remark label(bt:live_17:…)实例唯一——
|
||||
EOD 对账重启兜底归因(live_reconcile._own_orders_remarked)的前提。
|
||||
钉住 LiveConfig.load 吃这个键的 seam,上游升级若丢它必红。"""
|
||||
from bullet_trade.core.live_engine import LiveConfig
|
||||
cfg = LiveConfig.load({"strategy_name": "live_17"})
|
||||
assert cfg.strategy_name == "live_17"
|
||||
|
||||
Reference in New Issue
Block a user