跳过无业务分钟回调
This commit is contained in:
@@ -384,6 +384,10 @@ impl<C, R> BrokerSimulator<C, R> {
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub fn has_open_orders(&self) -> bool {
|
||||
!self.open_orders.borrow().is_empty()
|
||||
}
|
||||
}
|
||||
|
||||
impl<C, R> BrokerSimulator<C, R>
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
use std::collections::{BTreeMap, BTreeSet};
|
||||
|
||||
use chrono::{Datelike, Duration, NaiveDate, NaiveTime};
|
||||
use chrono::{Datelike, Duration, NaiveDate, NaiveTime, Timelike};
|
||||
use serde::Serialize;
|
||||
use thiserror::Error;
|
||||
|
||||
@@ -1086,6 +1086,10 @@ where
|
||||
views
|
||||
}
|
||||
|
||||
fn has_open_orders(&self) -> bool {
|
||||
self.broker.has_open_orders() || !self.futures_open_orders.is_empty()
|
||||
}
|
||||
|
||||
fn aggregate_initial_cash(&self) -> f64 {
|
||||
self.config.initial_cash
|
||||
+ self
|
||||
@@ -2486,6 +2490,20 @@ where
|
||||
!filter_by_subscription || self.subscriptions.contains("e.symbol)
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
let requires_minute_callbacks = self.strategy.requires_minute_callbacks();
|
||||
let has_minute_process_listeners = self.process_event_bus.has_listeners_for(&[
|
||||
ProcessEventKind::PreMinute,
|
||||
ProcessEventKind::Minute,
|
||||
ProcessEventKind::PostMinute,
|
||||
]);
|
||||
let minute_schedule_all_times = schedule_rules
|
||||
.iter()
|
||||
.any(|rule| rule.stage == ScheduleStage::Minute && rule.time_rule.is_none());
|
||||
let minute_schedule_minutes = schedule_rules
|
||||
.iter()
|
||||
.filter(|rule| rule.stage == ScheduleStage::Minute)
|
||||
.filter_map(|rule| rule.time_rule.as_ref()?.minute_of_day())
|
||||
.collect::<BTreeSet<_>>();
|
||||
let mut minute_cursor = 0usize;
|
||||
while minute_cursor < minute_quotes.len() {
|
||||
let minute_timestamp = minute_quotes[minute_cursor].timestamp;
|
||||
@@ -2497,6 +2515,17 @@ where
|
||||
minute_end += 1;
|
||||
}
|
||||
let minute_group = &minute_quotes[minute_cursor..minute_end];
|
||||
let schedule_candidate = minute_schedule_all_times
|
||||
|| minute_schedule_minutes
|
||||
.contains(&(minute_time.hour() * 60 + minute_time.minute()));
|
||||
if !requires_minute_callbacks
|
||||
&& !has_minute_process_listeners
|
||||
&& !schedule_candidate
|
||||
&& !self.has_open_orders()
|
||||
{
|
||||
minute_cursor = minute_end;
|
||||
continue;
|
||||
}
|
||||
let minute_open_orders = self.open_order_views();
|
||||
publish_phase_event(
|
||||
&mut self.strategy,
|
||||
@@ -2535,26 +2564,28 @@ where
|
||||
result.order_events.as_slice(),
|
||||
result.fills.as_slice(),
|
||||
)?;
|
||||
for quote in minute_group {
|
||||
minute_decision.merge_from(self.strategy.on_minute(
|
||||
&StrategyContext {
|
||||
execution_date,
|
||||
decision_date,
|
||||
decision_index,
|
||||
data: &self.data,
|
||||
portfolio: &portfolio,
|
||||
futures_account: self.futures_account.as_ref(),
|
||||
open_orders: &minute_open_orders,
|
||||
dynamic_universe: self.dynamic_universe.as_ref(),
|
||||
subscriptions: &self.subscriptions,
|
||||
process_events: &process_events,
|
||||
active_process_event: None,
|
||||
active_datetime: Some(minute_timestamp),
|
||||
order_events: result.order_events.as_slice(),
|
||||
fills: result.fills.as_slice(),
|
||||
},
|
||||
quote,
|
||||
)?);
|
||||
if requires_minute_callbacks {
|
||||
for quote in minute_group {
|
||||
minute_decision.merge_from(self.strategy.on_minute(
|
||||
&StrategyContext {
|
||||
execution_date,
|
||||
decision_date,
|
||||
decision_index,
|
||||
data: &self.data,
|
||||
portfolio: &portfolio,
|
||||
futures_account: self.futures_account.as_ref(),
|
||||
open_orders: &minute_open_orders,
|
||||
dynamic_universe: self.dynamic_universe.as_ref(),
|
||||
subscriptions: &self.subscriptions,
|
||||
process_events: &process_events,
|
||||
active_process_event: None,
|
||||
active_datetime: Some(minute_timestamp),
|
||||
order_events: result.order_events.as_slice(),
|
||||
fills: result.fills.as_slice(),
|
||||
},
|
||||
quote,
|
||||
)?);
|
||||
}
|
||||
}
|
||||
publish_phase_event(
|
||||
&mut self.strategy,
|
||||
|
||||
@@ -125,6 +125,15 @@ impl ProcessEventBus {
|
||||
loader.install_enabled(self, enabled_names)
|
||||
}
|
||||
|
||||
pub fn has_listeners_for(&self, kinds: &[ProcessEventKind]) -> bool {
|
||||
!self.any_listeners.is_empty()
|
||||
|| kinds.iter().any(|kind| {
|
||||
self.listeners
|
||||
.get(kind)
|
||||
.is_some_and(|listeners| !listeners.is_empty())
|
||||
})
|
||||
}
|
||||
|
||||
pub fn publish(&mut self, event: &ProcessEvent) {
|
||||
if let Some(listeners) = self.listeners.get_mut(&event.kind) {
|
||||
for listener in listeners {
|
||||
|
||||
@@ -10057,6 +10057,10 @@ impl Strategy for PlatformExprStrategy {
|
||||
self.config.initial_subscriptions.clone()
|
||||
}
|
||||
|
||||
fn requires_minute_callbacks(&self) -> bool {
|
||||
false
|
||||
}
|
||||
|
||||
fn schedule_rules(&self) -> Vec<ScheduleRule> {
|
||||
if self.config.explicit_action_stage != PlatformExplicitActionStage::Minute {
|
||||
return Vec::new();
|
||||
|
||||
@@ -22,6 +22,9 @@ pub trait Strategy {
|
||||
fn initial_subscriptions(&self) -> BTreeSet<String> {
|
||||
BTreeSet::new()
|
||||
}
|
||||
fn requires_minute_callbacks(&self) -> bool {
|
||||
true
|
||||
}
|
||||
fn management_fee(
|
||||
&mut self,
|
||||
_ctx: &StrategyContext<'_>,
|
||||
|
||||
Reference in New Issue
Block a user