From 191270c884eecfc1bad74d1d6e7f3be4728a9009 Mon Sep 17 00:00:00 2001 From: claude_dev Date: Tue, 25 Aug 2026 19:47:13 +0800 Subject: [PATCH] =?UTF-8?q?fix(trader):=20=E5=8D=96=E5=90=8E=E4=B9=B0?= =?UTF-8?q?=E7=8E=B0=E9=87=91=E7=AA=97=E5=8F=A3B=E4=BF=AE=E6=B3=95?= =?UTF-8?q?=E2=80=94=E2=80=94=E4=B8=8B=E5=8D=95=E8=BF=94=E5=9B=9E=E5=8D=B3?= =?UTF-8?q?=E6=97=B6=E5=BD=92=E5=9B=A0=E5=85=A5=E8=B4=A6,=E5=8F=B0?= =?UTF-8?q?=E8=B4=A6cash=E4=B8=8D=E5=86=8D=E7=AD=8960s=E8=BD=AE=E8=AF=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 2026-08-25 事故:small_cap同轮「全卖19只→马上全买20只」在两轮归因轮询间隙 读台账现金,卖出回款不可见→20笔买入全部目标0、全天空仓;momentum同型撞运 只入账首笔79k/6=13.2k缩水44%仓;value无卖后买序列满额(反证)。QMT无责 (0.5s filled/券商现金即时/下单线程同步见filled),gap=runner_live.py归因 poller 60s一轮才调ledger.apply_trade(唯一cash更新入口,DB落库时间戳恰差60s 铁证)。 修法(B,治本现金新鲜度): - LiveInstanceLedger.on_order_done钩子+notify_order_done(未注入/抛错静默, 绝不阻断下单;漏单由轮询兜底);apply_trade幂等判定整体移入锁内——钩子 (策略线程)与轮询(poller线程)并发同步同一笔成交时恰好一笔入账,防双计 - live_strategy._instance_order_wrappers:所有真实委托(bt_order/透传)返回后 _done()触发即时归因;决策层不下单的路径不触发 - runner_live:engine装配后注入on_order_done=_sync_instance_trades闭包; 60s轮询保留兜底(部分成交后续/异步路径) 测试+8:钩子三态(nop/触发/吞异常)+8线程同trade_id并发恰入账一次(竞态回归) +wrapper卖出/买入/透传触发+不下单不触发;portfolio 459绿+api 170绿 [vps] --- sanguo_portfolio/live_instance_ledger.py | 36 +++++++++++--- sanguo_portfolio/live_strategy.py | 21 +++++--- sanguo_portfolio/runner_live.py | 5 ++ tests/portfolio/test_live_instance_ledger.py | 51 ++++++++++++++++++++ tests/portfolio/test_live_instance_orders.py | 51 ++++++++++++++++++++ 5 files changed, 150 insertions(+), 14 deletions(-) diff --git a/sanguo_portfolio/live_instance_ledger.py b/sanguo_portfolio/live_instance_ledger.py index b5be2ed..b9a6c2d 100644 --- a/sanguo_portfolio/live_instance_ledger.py +++ b/sanguo_portfolio/live_instance_ledger.py @@ -22,7 +22,7 @@ from __future__ import annotations import logging import threading -from typing import Any, Dict, Iterable, Optional, Tuple +from typing import Any, Callable, Dict, Iterable, Optional, Tuple logger = logging.getLogger(__name__) @@ -59,6 +59,26 @@ class LiveInstanceLedger: self._lock = threading.Lock() # 有新成交未落 balance 快照 → 下个快照周期必写(节流档位见 runner_live) self.dirty = True + # B 修法(2026-08-25 卖后买现金窗口):下单返回后即时归因钩子, + # runner_live 注入 _sync_instance_trades 闭包;未注入(回测/影子/单测)=无操作 + self.on_order_done: Optional[Callable[[], None]] = None + + # ------------------ 即时归因入口(B修法) ------------------ + def notify_order_done(self) -> None: + """下单返回后立刻归因——台账 cash 秒级新鲜,不等 60s 归因轮询。 + + 2026-08-25 事故:small_cap 同轮「全卖→马上全买」在两轮轮询间隙读现金, + 19 笔卖出回款不可见 → 20 笔买入全部目标0、全天空仓(momentum 同型缩水 + 44% 仓)。钩子把引擎已见的成交即时喂进账本;未注入或抛错均静默—— + 漏掉的成交由归因轮询兜底,绝不阻断下单主流程。 + """ + hook = self.on_order_done + if hook is None: + return + try: + hook() + except Exception as e: # noqa: BLE001 + logger.warning("[instance-ledger] 即时归因失败,等60s轮询兜底: %s", e) # ------------------ 成交驱动 ------------------ def apply_trade( @@ -74,14 +94,16 @@ class LiveInstanceLedger: """应用一笔本实例成交;trade_id 重复返回 False(幂等)。 fee=None 时按费率估算;快照带实际佣金/印花税则传实际值。 + 幂等判定整体在锁内:即时归因钩子(策略线程)与归因轮询(poller 线程) + 并发同步同一笔成交时,恰好一笔入账(判定在锁外会双计现金)。 """ - if not trade_id or trade_id in self._seen_trade_ids: - return False - if price <= 0 or volume <= 0: - logger.warning("[instance-ledger] 非法成交跳过 %s %s x%s@%s", - trade_id, symbol, volume, price) - return False with self._lock: + if not trade_id or trade_id in self._seen_trade_ids: + return False + if price <= 0 or volume <= 0: + logger.warning("[instance-ledger] 非法成交跳过 %s %s x%s@%s", + trade_id, symbol, volume, price) + return False self._seen_trade_ids.add(trade_id) actual_fee = fee if (fee is not None and fee > 0) else \ estimate_fee(is_buy, price, volume) diff --git a/sanguo_portfolio/live_strategy.py b/sanguo_portfolio/live_strategy.py index 5b21840..9c28c66 100644 --- a/sanguo_portfolio/live_strategy.py +++ b/sanguo_portfolio/live_strategy.py @@ -140,6 +140,13 @@ def _instance_order_wrappers(ledger, bt_otv, bt_ov): return 0 return sell + def _done(order): + """B 修法(2026-08-25 卖后买现金窗口):真实委托返回后立刻即时归因—— + 引擎已见的成交即时进台账,cash 秒级新鲜,同轮「全卖→马上全买」不再 + 等不到卖出回款;无钩子(回测/影子/单测)为 no-op;不下单的路径不触发。""" + ledger.notify_order_done() + return order + def order_target_value(security, value, *args, **kwargs): own = _own(security) if value <= 0: @@ -148,42 +155,42 @@ def _instance_order_wrappers(ledger, bt_otv, bt_ov): return None logger.info("[instance-order] otv清仓 %s → 只卖自己 %d 股(目标值 %s)", security, sell, value) - return bt_order(security, -sell, *args, **kwargs) + return _done(bt_order(security, -sell, *args, **kwargs)) price = _price(security, own) if price <= 0: logger.warning("[instance-order] %s 无有效价格 → 透传引擎 target 语义", security) - return bt_otv(security, value, *args, **kwargs) + return _done(bt_otv(security, value, *args, **kwargs)) target = int(value / price) // lot * lot diff = target - int(own["amount"]) if diff >= lot: - return bt_order(security, diff // lot * lot, *args, **kwargs) + return _done(bt_order(security, diff // lot * lot, *args, **kwargs)) if diff < 0: sell = _sell_shares(security, -diff, own, "减仓") if sell <= 0: return None logger.info("[instance-order] otv减仓 %s → 卖 %d 股(目标 %d 现持 %d)", security, sell, target, int(own["amount"])) - return bt_order(security, -sell, *args, **kwargs) + return _done(bt_order(security, -sell, *args, **kwargs)) logger.info("[instance-order] otv %s 目标 %d ≈ 现持 %d → 不下单", security, target, int(own["amount"])) return None def order_value(security, value, *args, **kwargs): if value >= 0: - return bt_ov(security, value, *args, **kwargs) + return _done(bt_ov(security, value, *args, **kwargs)) own = _own(security) price = _price(security, own) if price <= 0: logger.warning("[instance-order] %s 无有效价格 → 负 order_value 透传引擎", security) - return bt_ov(security, value, *args, **kwargs) + return _done(bt_ov(security, value, *args, **kwargs)) want = int(abs(value) / price + 0.999) # 向上取整再由可卖量硬顶 sell = _sell_shares(security, want, own, "按价值卖出") if sell <= 0: return None logger.info("[instance-order] 负order_value %s → 卖 %d 股(价值 %s)", security, sell, value) - return bt_order(security, -sell, *args, **kwargs) + return _done(bt_order(security, -sell, *args, **kwargs)) return order_target_value, order_value diff --git a/sanguo_portfolio/runner_live.py b/sanguo_portfolio/runner_live.py index 0953c33..054db57 100644 --- a/sanguo_portfolio/runner_live.py +++ b/sanguo_portfolio/runner_live.py @@ -265,6 +265,11 @@ def run_live(provider_config: Dict[str, Any] | None = None) -> None: # 快照落库(supervisor 注入 db+account_id 时才开) if cfg["db"] and cfg["account_id"]: + # B 修法(2026-08-25 卖后买现金窗口):下单返回后即时归因——台账 cash + # 秒级新鲜,同轮「全卖→马上全买」的调仓立刻见到卖出回款;60s 归因 + # 轮询(_snapshot_loop)保留兜底(部分成交后续/异步路径)。 + ledger.on_order_done = lambda: _sync_instance_trades( + engine, ledger, cfg["db"], int(cfg["account_id"]), cfg["strategy"]) t = threading.Thread( target=_snapshot_loop, args=(engine, cfg["db"], int(cfg["account_id"]), ledger), diff --git a/tests/portfolio/test_live_instance_ledger.py b/tests/portfolio/test_live_instance_ledger.py index b2c13b9..10f643a 100644 --- a/tests/portfolio/test_live_instance_ledger.py +++ b/tests/portfolio/test_live_instance_ledger.py @@ -19,6 +19,57 @@ from sanguo_portfolio.live_instance_ledger import ( from sanguo_portfolio.runner_live import _snapshot_once, _sync_instance_trades +# ------------------ 即时归因钩子(B修法,2026-08-25 卖后买现金窗口) ------------------ +class TestOrderDoneHook: + """2026-08-25 事故:台账 cash 只被 60s 归因轮询更新,同轮「全卖→马上全买」 + 在轮询间隙读现金 → small_cap 20 笔买入全目标0 全天空仓、momentum 缩水44%仓。 + B 修法 = 下单返回后立刻归因(on_order_done 钩子),轮询降级为兜底。""" + + def test_notify_without_hook_is_noop(self): + led = LiveInstanceLedger() + led.notify_order_done() # 未注入钩子(回测/影子/单测)不抛不做事 + + def test_notify_invokes_registered_hook(self): + led = LiveInstanceLedger() + fired = [] + led.on_order_done = lambda: fired.append(1) + led.notify_order_done() + assert fired == [1] + + def test_notify_swallows_hook_exception(self): + """钩子失败绝不阻断下单主流程(漏单由 60s 轮询兜底)。""" + led = LiveInstanceLedger() + + def boom(): + raise RuntimeError("sync 失败") + + led.on_order_done = boom + led.notify_order_done() # 不抛 + + def test_concurrent_same_trade_id_credited_once(self): + """双线程竞态回归:即时归因钩子(策略线程)与轮询(poller线程)并发 + _sync 同一笔成交——幂等判定必须整体在锁内,否则现金双计。""" + import threading + + led = LiveInstanceLedger(initial_cash=100_000) + workers = [] + barrier = threading.Barrier(8) + + def worker(): + barrier.wait() + led.apply_trade(False, "000001.XSHE", 10.0, 1000, "same-tid", + "2026-08-25") + + for _ in range(8): + t = threading.Thread(target=worker) + t.start() + workers.append(t) + for t in workers: + t.join() + assert led.cash == pytest.approx( + 100_000 + 10_000 - estimate_fee(False, 10.0, 1000)) + + # ------------------ 账本算术 ------------------ class TestLedgerArithmetic: def test_buy_sell_with_estimated_fees(self): diff --git a/tests/portfolio/test_live_instance_orders.py b/tests/portfolio/test_live_instance_orders.py index 0650d46..68fb2bb 100644 --- a/tests/portfolio/test_live_instance_orders.py +++ b/tests/portfolio/test_live_instance_orders.py @@ -149,3 +149,54 @@ def test_price_falls_back_to_avg_cost(fakes, monkeypatch): otv, _ov = _instance_order_wrappers(led, None, None) otv("159915.XSHE", 6000) # 目标 3000 股 → 买 2900 assert orders == [("159915.XSHE", 2900)] and otv_calls == [] + + +# ------------------ 即时归因触发(B修法,2026-08-25 卖后买现金窗口) ------------------ +def _hooked(led): + fired = [] + led.on_order_done = lambda: fired.append(1) + return fired + + +def test_order_done_hook_fires_after_sell_order(fakes): + """真实卖单返回后立刻即时归因——台账现金不再等 60s 轮询才见卖出回款。""" + orders, _, _ = fakes + led = _ledger_with("513030.XSHG", 100) + fired = _hooked(led) + otv, _ov = _instance_order_wrappers(led, None, None) + otv("513030.XSHG", 0) + assert orders == [("513030.XSHG", -100)] + assert fired == [1] + + +def test_order_done_hook_fires_after_buy_order(fakes): + orders, _, _ = fakes + led = _ledger_with("512690.XSHG", 100, price=1.9) + fired = _hooked(led) + otv, _ov = _instance_order_wrappers(led, None, None) + otv("512690.XSHG", 5000) + assert orders == [("512690.XSHG", 2400)] + assert fired == [1] + + +def test_order_done_hook_fires_on_passthrough(fakes): + """透传引擎的路径(正 order_value)同样触发——真实委托都已出。""" + _, ov_calls, _ = fakes + led = _ledger_with("510300.XSHG", 0) + fired = _hooked(led) + fake_bt_ov = lambda sec, v, *a, **kw: ov_calls.append((sec, v)) or True + _otv, ov = _instance_order_wrappers(led, None, fake_bt_ov) + ov("510300.XSHG", 80000) + assert ov_calls == [("510300.XSHG", 80000)] + assert fired == [1] + + +def test_order_done_hook_silent_when_no_order(fakes): + """决策层就没下单(可卖0)→ 不触发即时归因(无新成交可喂)。""" + orders, _, _ = fakes + led = LiveInstanceLedger(initial_cash=500_000.0) + fired = _hooked(led) + otv, _ov = _instance_order_wrappers(led, None, None) + ret = otv("513030.XSHG", 0) + assert ret is None and orders == [] + assert fired == []