refactor(engine): 抽出step()单根推进(C-S3实走入口,task2基础)
run()循环体抽为step(bar_date,bars,prev_close,pending)→(pending,closes); run()改为调step。实走scheduler每日喂当日bar调step单步推进。 - 回放行为不变(test_engine原3用例pass) - 加test_engine_step_single_bar_advances: day1信号缓冲/day2撮合 - trader全量111/111通过
This commit is contained in:
+36
-29
@@ -54,41 +54,48 @@ class PaperEngine:
|
||||
self.interval = interval
|
||||
self.adjust = adjust
|
||||
|
||||
def step(self, bar_date, bars, prev_close, pending):
|
||||
"""单根 bar 推进(回放 run 循环调;C-S3 实走 scheduler 每日调)。
|
||||
|
||||
返回 (新 pending, 当根 closes)——实走每日喂当日 bar 调一次。
|
||||
"""
|
||||
self._bar_count = getattr(self, "_bar_count", 0) + 1
|
||||
self.account.unfreeze_all()
|
||||
for r in self.runners:
|
||||
r.unfreeze_all()
|
||||
# 1. 撮合上一根 pending(next_open,用当前 bar)
|
||||
if pending:
|
||||
for order, runner in pending:
|
||||
self._match(order, runner, bars, prev_close, bar_date)
|
||||
pending = []
|
||||
# 2. 喂策略 on_bar → 收新单
|
||||
for runner in self.runners:
|
||||
sym = runner.symbol
|
||||
if sym and sym in bars:
|
||||
runner.paper_cta_engine.on_bar(bars[sym])
|
||||
for order in runner.paper_cta_engine.pop_orders():
|
||||
if order.match_session == MatchSession.NEXT_OPEN:
|
||||
pending.append((order, runner))
|
||||
else: # current_close 当根撮合
|
||||
self._match(order, runner, bars, prev_close, bar_date)
|
||||
# 3. 盯市 + 入库
|
||||
closes = {s: bars[s].close_price for s in bars}
|
||||
self.account.mark_to_market(closes)
|
||||
save_daily_balance(
|
||||
self.db_path, self.account_id, str(bar_date),
|
||||
self.account.cash, self.account.market_value, self.account.equity,
|
||||
is_checkpoint=(self._bar_count % 500 == 0),
|
||||
)
|
||||
update_checkpoint(self.db_path, self.account_id, str(bar_date))
|
||||
return pending, closes
|
||||
|
||||
def run(self) -> None:
|
||||
prev_close: dict[str, float] = {}
|
||||
pending: list = [] # [(order, runner)] next_open 待下根撮合
|
||||
bar_count = 0
|
||||
for bar_date, bars in self.data_source.iter_bars(
|
||||
self.symbols, self.start, self.end, self.interval, self.adjust, None
|
||||
):
|
||||
bar_count += 1
|
||||
self.account.unfreeze_all()
|
||||
for r in self.runners:
|
||||
r.unfreeze_all()
|
||||
# 1. 撮合上一根 pending(next_open,用当前 bar)
|
||||
if pending:
|
||||
for order, runner in pending:
|
||||
self._match(order, runner, bars, prev_close, bar_date)
|
||||
pending = []
|
||||
# 2. 喂策略 on_bar → 收新单
|
||||
for runner in self.runners:
|
||||
sym = runner.symbol
|
||||
if sym and sym in bars:
|
||||
runner.paper_cta_engine.on_bar(bars[sym])
|
||||
for order in runner.paper_cta_engine.pop_orders():
|
||||
if order.match_session == MatchSession.NEXT_OPEN:
|
||||
pending.append((order, runner))
|
||||
else: # current_close 当根撮合
|
||||
self._match(order, runner, bars, prev_close, bar_date)
|
||||
# 3. 盯市 + 入库
|
||||
closes = {s: bars[s].close_price for s in bars}
|
||||
self.account.mark_to_market(closes)
|
||||
save_daily_balance(
|
||||
self.db_path, self.account_id, str(bar_date),
|
||||
self.account.cash, self.account.market_value, self.account.equity,
|
||||
is_checkpoint=(bar_count % 500 == 0),
|
||||
)
|
||||
update_checkpoint(self.db_path, self.account_id, str(bar_date))
|
||||
pending, closes = self.step(bar_date, bars, prev_close, pending)
|
||||
prev_close = closes
|
||||
|
||||
def _match(self, order, runner, bars, prev_close, bar_date) -> None:
|
||||
|
||||
Reference in New Issue
Block a user