fix(trader): 卖后买现金窗口B修法——下单返回即时归因入账,台账cash不再等60s轮询
CI/CD / test (push) Failing after 6s
CI/CD / nas-deploy (push) Has been skipped
CI/CD / nas-verify (push) Has been skipped

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]
This commit is contained in:
2026-08-25 19:47:13 +08:00
parent f438c6cd65
commit 191270c884
5 changed files with 150 additions and 14 deletions
+29 -7
View File
@@ -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)
+14 -7
View File
@@ -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
+5
View File
@@ -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),
@@ -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):
@@ -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 == []