支持任意交易阶段调度时间

This commit is contained in:
boris
2026-08-25 09:32:34 +08:00
parent 5ff8ddca92
commit c9ddff46dd
3 changed files with 125 additions and 42 deletions
+70 -12
View File
@@ -2043,7 +2043,7 @@ where
ProcessEventKind::BeforeTrading,
"before_trading",
)?;
let mut before_trading_decision = collect_scheduled_decisions(
let mut before_trading_decision = collect_scheduled_decisions_for_stage(
&mut self.strategy,
&scheduler,
execution_date,
@@ -2059,7 +2059,6 @@ where
&self.subscriptions,
&mut process_events,
&mut self.process_event_bus,
default_stage_time(ScheduleStage::BeforeTrading),
result.order_events.as_slice(),
result.fills.as_slice(),
)?;
@@ -2107,7 +2106,7 @@ where
ProcessEventKind::PreOpenAuction,
"open_auction:pre",
)?;
let mut auction_decision = collect_scheduled_decisions(
let mut auction_decision = collect_scheduled_decisions_for_stage(
&mut self.strategy,
&scheduler,
execution_date,
@@ -2123,7 +2122,6 @@ where
&self.subscriptions,
&mut process_events,
&mut self.process_event_bus,
default_stage_time(ScheduleStage::OpenAuction),
result.order_events.as_slice(),
result.fills.as_slice(),
)?;
@@ -2300,7 +2298,7 @@ where
})
.transpose()?
.unwrap_or_default();
decision.merge_from(collect_scheduled_decisions(
decision.merge_from(collect_scheduled_decisions_for_stage(
&mut self.strategy,
&scheduler,
execution_date,
@@ -2316,7 +2314,6 @@ where
&self.subscriptions,
&mut process_events,
&mut self.process_event_bus,
default_stage_time(ScheduleStage::OnDay),
result.order_events.as_slice(),
result.fills.as_slice(),
)?);
@@ -2355,7 +2352,7 @@ where
ProcessEventKind::PreBar,
"bar:pre",
)?;
decision.merge_from(collect_scheduled_decisions(
decision.merge_from(collect_scheduled_decisions_for_stage(
&mut self.strategy,
&scheduler,
execution_date,
@@ -2371,7 +2368,6 @@ where
&self.subscriptions,
&mut process_events,
&mut self.process_event_bus,
default_stage_time(ScheduleStage::Bar),
result.order_events.as_slice(),
result.fills.as_slice(),
)?);
@@ -2770,7 +2766,7 @@ where
ProcessEventKind::AfterTrading,
"after_trading",
)?;
let mut after_trading_decision = collect_scheduled_decisions(
let mut after_trading_decision = collect_scheduled_decisions_for_stage(
&mut self.strategy,
&scheduler,
execution_date,
@@ -2786,7 +2782,6 @@ where
&self.subscriptions,
&mut process_events,
&mut self.process_event_bus,
default_stage_time(ScheduleStage::AfterTrading),
result.order_events.as_slice(),
result.fills.as_slice(),
)?;
@@ -2899,7 +2894,7 @@ where
ProcessEventKind::Settlement,
"settlement",
)?;
let mut settlement_decision = collect_scheduled_decisions(
let mut settlement_decision = collect_scheduled_decisions_for_stage(
&mut self.strategy,
&scheduler,
execution_date,
@@ -2915,7 +2910,6 @@ where
&self.subscriptions,
&mut process_events,
&mut self.process_event_bus,
default_stage_time(ScheduleStage::Settlement),
result.order_events.as_slice(),
result.fills.as_slice(),
)?;
@@ -3845,6 +3839,70 @@ fn collect_scheduled_decisions<S: Strategy>(
Ok(combined)
}
fn collect_scheduled_decisions_for_stage<S: Strategy>(
strategy: &mut S,
scheduler: &Scheduler<'_>,
execution_date: NaiveDate,
stage: ScheduleStage,
rules: &[ScheduleRule],
decision_date: NaiveDate,
decision_index: usize,
data: &crate::data::DataSet,
portfolio: &PortfolioState,
futures_account: Option<&FuturesAccountState>,
open_orders: &[crate::strategy::OpenOrderView],
dynamic_universe: Option<&BTreeSet<String>>,
subscriptions: &BTreeSet<String>,
process_events: &mut Vec<ProcessEvent>,
process_event_bus: &mut ProcessEventBus,
order_events: &[OrderEvent],
fills: &[FillEvent],
) -> Result<crate::strategy::StrategyDecision, BacktestError> {
let mut times = BTreeSet::new();
for rule in rules.iter().filter(|rule| rule.stage == stage) {
let time = match rule.time_rule.as_ref() {
Some(crate::scheduler::ScheduleTimeRule::MinuteOfDay(value)) => {
let hour = value / 60;
let minute = value % 60;
Some(NaiveTime::from_hms_opt(hour, minute, 0).ok_or_else(|| {
BacktestError::Execution(format!(
"invalid schedule minute-of-day {} for rule {}",
value, rule.name
))
})?)
}
Some(crate::scheduler::ScheduleTimeRule::BeforeTrading) | None => {
default_stage_time(stage)
}
};
times.insert(time);
}
let mut combined = crate::strategy::StrategyDecision::default();
for time in times {
combined.merge_from(collect_scheduled_decisions(
strategy,
scheduler,
execution_date,
stage,
rules,
decision_date,
decision_index,
data,
portfolio,
futures_account,
open_orders,
dynamic_universe,
subscriptions,
process_events,
process_event_bus,
time,
order_events,
fills,
)?);
}
Ok(combined)
}
fn publish_phase_event<S: Strategy>(
strategy: &mut S,
process_event_bus: &mut ProcessEventBus,