1use qs_risk::{
4 ExposureFact, ExposureIntent, HaltCommand, IntentKind, PortfolioFacts, PortfolioSupervisor,
5 Verdict,
6};
7
8use super::*;
9use crate::artifacts::RecordedFill;
10use crate::strategy::{
11 ConfiguredInstance, DirectPortfolioInstance, INSTANCE_POSITION_TAG, MAX_PORTFOLIO_INSTANCES,
12 MixedPortfolioBacktestResult, MixedPortfolioReplayError, PortfolioBacktestResult,
13 PortfolioInstanceOutput, PortfolioReplayError, SupervisorEvent, SupervisorHaltAction,
14 SupervisorOutput,
15};
16
17pub(super) struct PortfolioReplayState {
19 instance_ids: Vec<String>,
20 supervisor: Option<PortfolioSupervisor>,
21 reservations: BTreeMap<String, ExposureFact>,
23 disposed: BTreeSet<String>,
25 seen_dispositions: usize,
26 action_instances: BTreeMap<String, usize>,
28 events: Vec<SupervisorEvent>,
29 halt_actions: Vec<SupervisorHaltAction>,
30 unmarked_boundaries: u64,
31 next_halt_action: u64,
32}
33
34impl PortfolioReplayState {
35 fn new(instance_ids: Vec<String>, supervisor: Option<PortfolioSupervisor>) -> Self {
36 Self {
37 instance_ids,
38 supervisor,
39 reservations: BTreeMap::new(),
40 disposed: BTreeSet::new(),
41 seen_dispositions: 0,
42 action_instances: BTreeMap::new(),
43 events: Vec::new(),
44 halt_actions: Vec::new(),
45 unmarked_boundaries: 0,
46 next_halt_action: 0,
47 }
48 }
49
50 pub(super) fn begin_batch(&mut self, ts: NaiveDateTime, balance: f64) {
51 if let Some(supervisor) = self.supervisor.as_mut() {
52 supervisor.begin(ts, balance);
53 }
54 }
55
56 pub(super) fn position_tags(
58 &self,
59 fills: &[RecordedFill],
60 ) -> BTreeMap<String, BTreeMap<String, String>> {
61 let mut tags = BTreeMap::new();
62 for fill in fills {
63 let Some(instance) = fill
64 .action_id
65 .as_ref()
66 .and_then(|action_id| self.action_instances.get(action_id))
67 else {
68 continue;
69 };
70 tags.entry(fill.position_id.clone()).or_insert_with(|| {
71 BTreeMap::from([(
72 INSTANCE_POSITION_TAG.to_owned(),
73 self.instance_ids[*instance].clone(),
74 )])
75 });
76 }
77 tags
78 }
79
80 fn into_output(self) -> Option<SupervisorOutput> {
81 let supervisor = self.supervisor?;
82 Some(SupervisorOutput {
83 events: self.events,
84 halt_actions: self.halt_actions,
85 halts: supervisor.finish(),
86 unmarked_boundaries: self.unmarked_boundaries,
87 })
88 }
89}
90
91struct PortfolioReplayHook<'a> {
93 drivers: Vec<ConfiguredStrategyReplayDriver<'a>>,
94 declared: Vec<BTreeMap<String, BTreeSet<u64>>>,
96 feed_from: Vec<Option<NaiveDateTime>>,
98 state: PortfolioReplayState,
99 failed: Option<usize>,
100 seen_input: (bool, bool),
102 mixed_input: Option<NaiveDateTime>,
104}
105
106impl PortfolioReplayHook<'_> {
107 fn reads(&self, instance: usize, event: &FeedEvent) -> bool {
108 if self.feed_from[instance].is_some_and(|from| event.event.ts() < from) {
109 return false;
110 }
111 let Some(durations) = self.declared[instance].get(event.event.symbol()) else {
112 return false;
113 };
114 match &event.event {
115 MarketEvent::Tick { .. } => true,
116 MarketEvent::Bar {
117 timeframe_seconds, ..
118 } => timeframe_seconds.is_none_or(|seconds| durations.contains(&seconds)),
119 }
120 }
121}
122
123impl FutureReplayHook for PortfolioReplayHook<'_> {
124 fn is_active(&self) -> bool {
125 true
126 }
127
128 fn output_ready(&self) -> bool {
129 self.drivers.iter().any(FutureReplayHook::output_ready)
130 }
131
132 fn preflight_primary_events(&mut self, events: &[FeedEvent]) -> bool {
134 for event in events {
135 match event.event {
136 MarketEvent::Tick { .. } => self.seen_input.0 = true,
137 MarketEvent::Bar { .. } => self.seen_input.1 = true,
138 }
139 if self.seen_input == (true, true) {
140 self.mixed_input = Some(event.event.ts());
141 return false;
142 }
143 }
144 for (index, driver) in self.drivers.iter_mut().enumerate() {
145 if !driver.preflight_primary_events(events) {
146 self.failed = Some(index);
147 return false;
148 }
149 }
150 true
151 }
152
153 fn reject_generated_configuration(&mut self, instance: Option<usize>, reason: String) {
154 let index = instance.unwrap_or(0);
155 self.failed = Some(index);
156 self.drivers[index].reject_generated_configuration(None, reason);
157 }
158
159 fn observes_position_economics(&self) -> bool {
160 true
161 }
162
163 fn reads_completed_bars_only(&self) -> bool {
164 true
165 }
166
167 fn bar_execution_timeframes(&self) -> BTreeMap<String, u64> {
168 let mut shortest = BTreeMap::<String, u64>::new();
169 for driver in &self.drivers {
170 for (symbol, seconds) in driver.bar_execution_timeframes() {
171 shortest
172 .entry(symbol)
173 .and_modify(|current| *current = (*current).min(seconds))
174 .or_insert(seconds);
175 }
176 }
177 shortest
178 }
179
180 fn portfolio_state(&mut self) -> Option<&mut PortfolioReplayState> {
181 Some(&mut self.state)
182 }
183
184 fn on_boundary(
185 &mut self,
186 batch: &TimestampBatch,
187 engine: &TradeEngine,
188 lifecycle: &LifecycleLedger,
189 positions: &BoundaryPositionFacts<'_>,
190 pending_effects: &mut Vec<FutureEffect>,
191 pending_events: &mut Vec<StrategyFeedbackEvent>,
192 ) -> Option<Vec<ScheduledSignal>> {
193 let mut scheduled = Vec::new();
194 for index in 0..self.drivers.len() {
195 let own_batch = TimestampBatch {
196 ts: batch.ts,
197 events: batch
198 .events
199 .iter()
200 .filter(|event| self.reads(index, event))
201 .cloned()
202 .collect(),
203 };
204 let mut effects = pending_effects.clone();
206 let mut events = pending_events.clone();
207 let Some(output) = self.drivers[index].on_boundary(
208 &own_batch,
209 engine,
210 lifecycle,
211 positions,
212 &mut effects,
213 &mut events,
214 ) else {
215 self.failed = Some(index);
216 return None;
217 };
218 scheduled.extend(output.into_iter().map(|mut signal| {
219 signal.instance = Some(index);
220 signal
221 }));
222 }
223 pending_effects.clear();
224 pending_events.clear();
225 Some(scheduled)
226 }
227
228 fn on_final_committed(
229 &mut self,
230 pending_effects: &mut Vec<FutureEffect>,
231 pending_events: &mut Vec<StrategyFeedbackEvent>,
232 ) -> bool {
233 for index in 0..self.drivers.len() {
234 let mut effects = pending_effects.clone();
235 let mut events = pending_events.clone();
236 if !self.drivers[index].on_final_committed(&mut effects, &mut events) {
237 self.failed = Some(index);
238 return false;
239 }
240 }
241 pending_effects.clear();
242 pending_events.clear();
243 true
244 }
245}
246
247enum MixedDriverFailure {
248 Configured(usize),
249 Direct(usize),
250}
251
252struct MixedPortfolioReplayHook<'a, E> {
253 configured: Vec<ConfiguredStrategyReplayDriver<'a>>,
254 direct: Vec<StrategyReplayDriver<'a, dyn HistoricalStrategy<Error = E> + Send>>,
255 configured_declared: Vec<BTreeMap<String, BTreeSet<u64>>>,
256 direct_declared: Vec<BTreeMap<String, BTreeSet<u64>>>,
257 configured_feed_from: Vec<Option<NaiveDateTime>>,
258 direct_feed_from: Vec<Option<NaiveDateTime>>,
259 state: PortfolioReplayState,
260 failed: Option<MixedDriverFailure>,
261 seen_input: (bool, bool),
262 mixed_input: Option<NaiveDateTime>,
263}
264
265impl<E> MixedPortfolioReplayHook<'_, E> {
266 fn reads(
267 declared: &[BTreeMap<String, BTreeSet<u64>>],
268 feed_from: &[Option<NaiveDateTime>],
269 instance: usize,
270 event: &FeedEvent,
271 ) -> bool {
272 if feed_from[instance].is_some_and(|from| event.available_at() < from) {
273 return false;
274 }
275 let Some(durations) = declared[instance].get(event.event.symbol()) else {
276 return false;
277 };
278 match &event.event {
279 MarketEvent::Tick { .. } => true,
280 MarketEvent::Bar {
281 timeframe_seconds, ..
282 } => timeframe_seconds.is_none_or(|seconds| durations.contains(&seconds)),
283 }
284 }
285
286 #[allow(clippy::too_many_arguments)]
287 fn drive_configured(
288 &mut self,
289 batch: &TimestampBatch,
290 engine: &TradeEngine,
291 lifecycle: &LifecycleLedger,
292 positions: &BoundaryPositionFacts<'_>,
293 pending_effects: &[FutureEffect],
294 pending_events: &[StrategyFeedbackEvent],
295 scheduled: &mut Vec<ScheduledSignal>,
296 ) -> bool {
297 for index in 0..self.configured.len() {
298 let own_batch = TimestampBatch {
299 ts: batch.ts,
300 events: batch
301 .events
302 .iter()
303 .filter(|event| {
304 Self::reads(
305 &self.configured_declared,
306 &self.configured_feed_from,
307 index,
308 event,
309 )
310 })
311 .cloned()
312 .collect(),
313 };
314 let mut effects = pending_effects.to_vec();
315 let mut events = pending_events.to_vec();
316 let Some(output) = self.configured[index].on_boundary(
317 &own_batch,
318 engine,
319 lifecycle,
320 positions,
321 &mut effects,
322 &mut events,
323 ) else {
324 self.failed = Some(MixedDriverFailure::Configured(index));
325 return false;
326 };
327 scheduled.extend(output.into_iter().map(|mut signal| {
328 signal.instance = Some(index);
329 signal
330 }));
331 }
332 true
333 }
334
335 #[allow(clippy::too_many_arguments)]
336 fn drive_direct(
337 &mut self,
338 batch: &TimestampBatch,
339 engine: &TradeEngine,
340 lifecycle: &LifecycleLedger,
341 positions: &BoundaryPositionFacts<'_>,
342 pending_effects: &[FutureEffect],
343 pending_events: &[StrategyFeedbackEvent],
344 scheduled: &mut Vec<ScheduledSignal>,
345 ) -> bool {
346 let offset = self.configured.len();
347 for index in 0..self.direct.len() {
348 let own_batch = TimestampBatch {
349 ts: batch.ts,
350 events: batch
351 .events
352 .iter()
353 .filter(|event| {
354 Self::reads(&self.direct_declared, &self.direct_feed_from, index, event)
355 })
356 .cloned()
357 .collect(),
358 };
359 let mut effects = pending_effects.to_vec();
360 let mut events = pending_events.to_vec();
361 let Some(output) = self.direct[index].on_boundary(
362 &own_batch,
363 engine,
364 lifecycle,
365 positions,
366 &mut effects,
367 &mut events,
368 ) else {
369 self.failed = Some(MixedDriverFailure::Direct(index));
370 return false;
371 };
372 scheduled.extend(output.into_iter().map(|mut signal| {
373 signal.instance = Some(offset + index);
374 signal
375 }));
376 }
377 true
378 }
379}
380
381impl<E> FutureReplayHook for MixedPortfolioReplayHook<'_, E> {
382 fn is_active(&self) -> bool {
383 true
384 }
385
386 fn output_ready(&self) -> bool {
387 self.configured.iter().any(FutureReplayHook::output_ready)
388 || self.direct.iter().any(FutureReplayHook::output_ready)
389 }
390
391 fn preflight_primary_events(&mut self, events: &[FeedEvent]) -> bool {
392 for event in events {
393 match event.event {
394 MarketEvent::Tick { .. } => self.seen_input.0 = true,
395 MarketEvent::Bar { .. } => self.seen_input.1 = true,
396 }
397 if self.seen_input == (true, true) {
398 self.mixed_input = Some(event.available_at());
399 return false;
400 }
401 }
402 for (index, driver) in self.configured.iter_mut().enumerate() {
403 if !driver.preflight_primary_events(events) {
404 self.failed = Some(MixedDriverFailure::Configured(index));
405 return false;
406 }
407 }
408 for (index, driver) in self.direct.iter_mut().enumerate() {
409 if !driver.preflight_primary_events(events) {
410 self.failed = Some(MixedDriverFailure::Direct(index));
411 return false;
412 }
413 }
414 true
415 }
416
417 fn reject_generated_configuration(&mut self, instance: Option<usize>, reason: String) {
418 let index = instance.unwrap_or(0);
419 if index < self.configured.len() {
420 self.failed = Some(MixedDriverFailure::Configured(index));
421 self.configured[index].reject_generated_configuration(None, reason);
422 } else {
423 let direct = index - self.configured.len();
424 self.failed = Some(MixedDriverFailure::Direct(direct));
425 self.direct[direct].reject_generated_configuration(None, reason);
426 }
427 }
428
429 fn observes_position_economics(&self) -> bool {
430 true
431 }
432
433 fn reads_completed_bars_only(&self) -> bool {
434 true
435 }
436
437 fn retains_post_bar_boundary(&self) -> bool {
438 true
439 }
440
441 fn bar_execution_timeframes(&self) -> BTreeMap<String, u64> {
442 let mut shortest = BTreeMap::new();
443 for driver in &self.configured {
444 for (symbol, seconds) in driver.bar_execution_timeframes() {
445 shortest
446 .entry(symbol)
447 .and_modify(|current: &mut u64| *current = (*current).min(seconds))
448 .or_insert(seconds);
449 }
450 }
451 for driver in &self.direct {
452 for (symbol, seconds) in driver.bar_execution_timeframes() {
453 shortest
454 .entry(symbol)
455 .and_modify(|current| *current = (*current).min(seconds))
456 .or_insert(seconds);
457 }
458 }
459 shortest
460 }
461
462 fn portfolio_state(&mut self) -> Option<&mut PortfolioReplayState> {
463 Some(&mut self.state)
464 }
465
466 fn on_pre_bar_boundary(
467 &mut self,
468 batch: &TimestampBatch,
469 engine: &TradeEngine,
470 lifecycle: &LifecycleLedger,
471 positions: &BoundaryPositionFacts<'_>,
472 pending_effects: &mut Vec<FutureEffect>,
473 pending_events: &mut Vec<StrategyFeedbackEvent>,
474 ) -> Option<Vec<ScheduledSignal>> {
475 let mut scheduled = Vec::new();
476 self.drive_configured(
477 batch,
478 engine,
479 lifecycle,
480 positions,
481 pending_effects,
482 pending_events,
483 &mut scheduled,
484 )
485 .then_some(scheduled)
486 }
487
488 fn on_boundary(
489 &mut self,
490 batch: &TimestampBatch,
491 engine: &TradeEngine,
492 lifecycle: &LifecycleLedger,
493 positions: &BoundaryPositionFacts<'_>,
494 pending_effects: &mut Vec<FutureEffect>,
495 pending_events: &mut Vec<StrategyFeedbackEvent>,
496 ) -> Option<Vec<ScheduledSignal>> {
497 let mut scheduled = Vec::new();
498 let bar_batch = batch
499 .events
500 .iter()
501 .any(|event| matches!(event.event, MarketEvent::Bar { .. }));
502 if !bar_batch
503 && !self.drive_configured(
504 batch,
505 engine,
506 lifecycle,
507 positions,
508 pending_effects,
509 pending_events,
510 &mut scheduled,
511 )
512 {
513 return None;
514 }
515 if !self.drive_direct(
516 batch,
517 engine,
518 lifecycle,
519 positions,
520 pending_effects,
521 pending_events,
522 &mut scheduled,
523 ) {
524 return None;
525 }
526 pending_effects.clear();
527 pending_events.clear();
528 Some(scheduled)
529 }
530
531 fn on_final_committed(
532 &mut self,
533 pending_effects: &mut Vec<FutureEffect>,
534 pending_events: &mut Vec<StrategyFeedbackEvent>,
535 ) -> bool {
536 for (index, driver) in self.configured.iter_mut().enumerate() {
537 let mut effects = pending_effects.clone();
538 let mut events = pending_events.clone();
539 if !driver.on_final_committed(&mut effects, &mut events) {
540 self.failed = Some(MixedDriverFailure::Configured(index));
541 return false;
542 }
543 }
544 for (index, driver) in self.direct.iter_mut().enumerate() {
545 let mut effects = pending_effects.clone();
546 let mut events = pending_events.clone();
547 if !driver.on_final_committed(&mut effects, &mut events) {
548 self.failed = Some(MixedDriverFailure::Direct(index));
549 return false;
550 }
551 }
552 pending_effects.clear();
553 pending_events.clear();
554 true
555 }
556}
557
558impl BacktestRunner {
559 #[allow(clippy::too_many_arguments)]
563 pub(super) fn supervise_generated<H: FutureReplayHook>(
564 &mut self,
565 hook: &mut H,
566 batch_ts: NaiveDateTime,
567 generated: Vec<ScheduledSignal>,
568 decided_before_quotes: bool,
569 drawdown_fraction: Option<f64>,
570 scheduled: &mut VecDeque<ScheduledSignal>,
571 queued: &mut VecDeque<QueuedAction>,
572 lifecycle: &mut LifecycleLedger,
573 future_executor: &FutureExecutor,
574 ) -> Vec<ScheduledSignal> {
575 let Some(state) = hook.portfolio_state() else {
576 return generated;
577 };
578 for signal in &generated {
579 if let (Some(instance), Some(action_id)) = (signal.instance, signal.action_id.as_ref())
580 {
581 state.action_instances.insert(action_id.clone(), instance);
582 }
583 }
584 if state.supervisor.is_none() {
585 return generated;
586 }
587 for disposition in &lifecycle.as_slice()[state.seen_dispositions..] {
588 state.disposed.insert(disposition.action_id.clone());
589 }
590 state.seen_dispositions = lifecycle.len();
591 if drawdown_fraction.is_none() {
592 state.unmarked_boundaries += 1;
593 }
594
595 let open = self
596 .engine
597 .open_positions()
598 .into_iter()
599 .map(|position| ExposureFact {
600 symbol: position.data.symbol.clone(),
601 side: position.data.side,
602 risk: future_executor.open_initial_risk(&position.data.id),
603 })
604 .collect::<Vec<_>>();
605 let pending = self
606 .engine
607 .pending_positions()
608 .into_iter()
609 .map(|position| ExposureFact {
610 symbol: position.data.symbol.clone(),
611 side: position.data.side,
612 risk: future_executor
613 .pending_metadata(&position.data.id)
614 .and_then(|(action_id, ..)| state.reservations.get(&action_id))
615 .and_then(|reserved| reserved.risk),
616 })
617 .collect::<Vec<_>>();
618 let mut reserved = state
619 .reservations
620 .iter()
621 .filter(|(action_id, _)| !state.disposed.contains(*action_id))
622 .map(|(_, fact)| fact.clone())
623 .collect::<Vec<_>>();
624 let balance = future_executor.balance();
625 let supervisor = state.supervisor.as_mut().expect("supervisor checked above");
626 let day_realized_r = supervisor.day_start().map_or(0.0, |start| {
627 future_executor
628 .completed_positions
629 .iter()
630 .filter(|position| position.close_ts >= start)
631 .filter_map(|position| position.realized_r)
632 .sum()
633 });
634 let halt_commands = supervisor.on_boundary(&PortfolioFacts {
635 now: batch_ts,
636 balance,
637 drawdown_fraction,
638 day_realized_r,
639 open: &open,
640 pending: &pending,
641 reserved: &reserved,
642 });
643 let halt_began = !halt_commands.is_empty();
644
645 let mut approved = Vec::with_capacity(generated.len());
646 let mut rejections = Vec::new();
647 for signal in generated {
648 let (symbol, side, kind, requested_risk) = match &signal.signal {
649 RawSignal::Entry {
650 symbol,
651 side,
652 risk_multiplier,
653 ..
654 } => (
655 symbol.clone(),
656 *side,
657 IntentKind::Entry,
658 requested_account_risk(self.config.sizing.as_ref(), balance, *risk_multiplier),
659 ),
660 RawSignal::ScaleIn { .. } => {
661 let Some((symbol, side)) = self
662 .resolve_future_actions(&signal.signal)
663 .into_iter()
664 .find_map(|action| match action {
665 Action::ScaleIn { position_id, .. } => self
666 .engine
667 .get_position(&position_id)
668 .map(|position| (position.data.symbol.clone(), position.data.side)),
669 _ => None,
670 })
671 else {
672 approved.push(signal);
673 continue;
674 };
675 (symbol, side, IntentKind::ScaleIn, None)
676 }
677 _ => {
678 approved.push(signal);
679 continue;
680 }
681 };
682 let verdict = supervisor.review(
683 &PortfolioFacts {
684 now: batch_ts,
685 balance,
686 drawdown_fraction,
687 day_realized_r,
688 open: &open,
689 pending: &pending,
690 reserved: &reserved,
691 },
692 &ExposureIntent {
693 symbol: &symbol,
694 side,
695 kind,
696 requested_risk,
697 },
698 );
699 let action_id = signal.resolved_action_id();
700 let instance_id = signal
701 .instance
702 .map(|instance| state.instance_ids[instance].clone())
703 .unwrap_or_default();
704 state.events.push(SupervisorEvent {
705 ts: batch_ts,
706 instance_id,
707 action_id: action_id.clone(),
708 kind,
709 symbol: symbol.clone(),
710 requested_risk,
711 verdict: verdict.clone(),
712 });
713 match verdict {
714 Verdict::Approve => {
715 if kind == IntentKind::Entry {
716 let fact = ExposureFact {
717 symbol,
718 side,
719 risk: requested_risk,
720 };
721 state.reservations.insert(action_id, fact.clone());
722 reserved.push(fact);
723 }
724 approved.push(signal);
725 }
726 Verdict::Reject { policy, reason } => {
727 let mut disposition =
728 ActionDisposition::rejected(action_id, format!("{policy}: {reason}"));
729 disposition.action_kind = Some(raw_signal_kind(&signal.signal).to_owned());
730 disposition.signal_ts = Some(signal.signal_ts);
731 disposition.effective_ts = Some(signal.effective_ts);
732 rejections.push(disposition);
733 }
734 }
735 }
736 for command in halt_commands {
737 let (signal, name) = match command {
738 HaltCommand::CancelAllPending => (
739 RawSignal::CancelAllPending { ts: batch_ts },
740 "cancel_all_pending",
741 ),
742 HaltCommand::CloseAll => (RawSignal::CloseAll { ts: batch_ts }, "close_all"),
743 };
744 let action_id = format!("supervisor:{:08}:{name}", state.next_halt_action);
745 state.next_halt_action += 1;
746 state.halt_actions.push(SupervisorHaltAction {
747 ts: batch_ts,
748 action_id: action_id.clone(),
749 command,
750 });
751 approved.push(
752 ScheduledSignal::new(0, batch_ts, batch_ts, signal, !decided_before_quotes)
753 .with_action_base(action_id),
754 );
755 }
756 if halt_began {
757 let policy = supervisor
758 .intervals()
759 .last()
760 .map_or_else(|| "halt".to_owned(), |interval| interval.policy.clone());
761 let reason = format!(
762 "{policy}: new exposure is halted and the request had not reached the market"
763 );
764 let mut remaining = VecDeque::with_capacity(queued.len());
765 for action in queued.drain(..) {
766 if action.entry_signal.is_some()
767 || matches!(action.action, Action::Open { .. } | Action::ScaleIn { .. })
768 {
769 let mut disposition =
770 ActionDisposition::rejected(action.action_id, reason.clone());
771 disposition.action_kind = Some(action.action_kind);
772 disposition.signal_ts = Some(action.signal_ts);
773 disposition.effective_ts = Some(action.effective_ts);
774 rejections.push(disposition);
775 } else {
776 remaining.push_back(action);
777 }
778 }
779 *queued = remaining;
780 let mut remaining = VecDeque::with_capacity(scheduled.len());
781 for signal in scheduled.drain(..) {
782 if matches!(
783 signal.signal,
784 RawSignal::Entry { .. } | RawSignal::ScaleIn { .. }
785 ) {
786 let mut disposition =
787 ActionDisposition::rejected(signal.resolved_action_id(), reason.clone());
788 disposition.action_kind = Some(raw_signal_kind(&signal.signal).to_owned());
789 disposition.signal_ts = Some(signal.signal_ts);
790 disposition.effective_ts = Some(signal.effective_ts);
791 rejections.push(disposition);
792 } else {
793 remaining.push_back(signal);
794 }
795 }
796 *scheduled = remaining;
797 }
798 for disposition in rejections {
799 self.record_disposition(lifecycle, disposition);
800 }
801 approved
802 }
803
804 pub fn run_portfolio_future<F>(
806 self,
807 source_feed: &mut F,
808 instances: Vec<ConfiguredInstance>,
809 supervisor: Option<PortfolioSupervisor>,
810 retention: StrategyRetentionLimits,
811 ) -> Result<PortfolioBacktestResult, PortfolioReplayError<Infallible>>
812 where
813 F: DataFeed,
814 {
815 let mut ordered_events = Vec::new();
816 let mut source_last_ts = BTreeMap::<String, NaiveDateTime>::new();
817 while let Some(batch) = source_feed.next_batch() {
818 for event in batch.events {
819 let symbol = event.event.symbol().to_owned();
820 let timestamp = event.event.ts();
821 if source_last_ts
822 .get(&symbol)
823 .is_some_and(|previous| *previous > timestamp)
824 {
825 continue;
826 }
827 source_last_ts.insert(symbol, timestamp);
828 ordered_events.push(event);
829 }
830 }
831 ordered_events.sort_by_key(FeedEvent::ordering_key);
832 let primary_eod = ordered_events
833 .iter()
834 .filter(|event| event.metadata.roles.primary)
835 .filter_map(|event| event.event.to_valid_quote())
836 .map(|quote| quote.ts)
837 .max();
838 let mut ordered_feed = crate::data_feed::VecFeed::from_feed_events(ordered_events);
839 let mut feed = DataFeedBatchAdapter {
840 feed: &mut ordered_feed,
841 };
842 self.run_portfolio_future_streaming_controlled(
843 &mut feed,
844 primary_eod,
845 instances,
846 supervisor,
847 retention,
848 || false,
849 |_| {},
850 )
851 }
852
853 #[allow(clippy::too_many_arguments)]
855 pub fn run_mixed_portfolio_future<F, E>(
856 self,
857 source_feed: &mut F,
858 configured: Vec<ConfiguredInstance>,
859 direct: Vec<DirectPortfolioInstance<E>>,
860 supervisor: Option<PortfolioSupervisor>,
861 retention: StrategyRetentionLimits,
862 ) -> Result<MixedPortfolioBacktestResult, MixedPortfolioReplayError<Infallible>>
863 where
864 F: DataFeed,
865 E: std::fmt::Display,
866 {
867 let mut events = Vec::new();
868 while let Some(batch) = source_feed.next_batch() {
869 events.extend(batch.events);
870 }
871 events.sort_by_key(FeedEvent::ordering_key);
872 let primary_eod = events
873 .iter()
874 .filter(|event| event.metadata.roles.primary)
875 .map(FeedEvent::available_at)
876 .max();
877 let mut feed = crate::data_feed::VecFeed::from_feed_events(events);
878 let mut feed = DataFeedBatchAdapter { feed: &mut feed };
879 self.run_mixed_portfolio_future_streaming_controlled(
880 &mut feed,
881 primary_eod,
882 configured,
883 direct,
884 supervisor,
885 retention,
886 || false,
887 |_| {},
888 )
889 }
890
891 #[allow(clippy::too_many_arguments)]
892 pub fn run_mixed_portfolio_future_streaming_controlled<F, E, C, P>(
893 mut self,
894 feed: &mut F,
895 primary_eod: Option<NaiveDateTime>,
896 configured: Vec<ConfiguredInstance>,
897 direct: Vec<DirectPortfolioInstance<E>>,
898 supervisor: Option<PortfolioSupervisor>,
899 retention: StrategyRetentionLimits,
900 mut is_cancelled: C,
901 mut on_progress: P,
902 ) -> Result<MixedPortfolioBacktestResult, MixedPortfolioReplayError<F::Error>>
903 where
904 F: FallibleBatchFeed,
905 E: std::fmt::Display,
906 C: FnMut() -> bool,
907 P: FnMut(ReplayProgress),
908 {
909 let total = configured.len().checked_add(direct.len()).ok_or_else(|| {
910 MixedPortfolioReplayError::Input("mixed instance count overflowed".into())
911 })?;
912 if total == 0 {
913 return Err(MixedPortfolioReplayError::NoInstances);
914 }
915 if total > MAX_PORTFOLIO_INSTANCES {
916 return Err(MixedPortfolioReplayError::TooManyInstances(total));
917 }
918 if self.entry_profiles.is_some() {
919 return Err(MixedPortfolioReplayError::Input(
920 "mixed portfolio instances carry their own entry profiles".into(),
921 ));
922 }
923 if self.config.run_tags.contains_key(INSTANCE_POSITION_TAG) {
924 return Err(MixedPortfolioReplayError::Input(format!(
925 "run tag '{INSTANCE_POSITION_TAG}' is owned by portfolio replay"
926 )));
927 }
928 if let Some(supervisor) = supervisor.as_ref()
929 && supervisor.caps_group_risk()
930 && !self.config.sizing.as_ref().is_some_and(is_monetary_sizing)
931 {
932 return Err(MixedPortfolioReplayError::Input(
933 "a group risk cap needs a monetary sizing policy".into(),
934 ));
935 }
936 let future = self.future_config.clone().unwrap_or_default();
937 self.future_config = Some(future.clone());
938 validate_replay_config(&self.config, Some(&future), &[])
939 .map_err(MixedPortfolioReplayError::Input)?;
940
941 let mut identities = BTreeSet::new();
942 let mut adapters = Vec::with_capacity(configured.len());
943 let mut configured_analysis = Vec::with_capacity(configured.len());
944 let mut configured_series = Vec::with_capacity(configured.len());
945 let mut configured_feed_from = Vec::with_capacity(configured.len());
946 let mut profiles = Vec::with_capacity(total);
947 let mut instance_ids = Vec::with_capacity(total);
948 for instance in configured {
949 let instance_id = instance.instance_id().to_owned();
950 if !identities.insert(instance_id.clone()) {
951 return Err(MixedPortfolioReplayError::DuplicateInstanceIdentity { instance_id });
952 }
953 instance
954 .adapter
955 .preflight_entry_profiles(&instance.entry_profiles)
956 .map_err(|error| MixedPortfolioReplayError::Instance {
957 instance_id: instance_id.clone(),
958 reason: error.to_string(),
959 })?;
960 instance
961 .adapter
962 .preflight(retention, self.strategy_research_limits)
963 .map_err(|error| MixedPortfolioReplayError::Instance {
964 instance_id: instance_id.clone(),
965 reason: error.to_string(),
966 })?;
967 let specs = instance.adapter.series_specs().cloned().collect::<Vec<_>>();
968 crate::strategy::replay::validate_series_specs(instance.adapter.requirements(), &specs)
969 .map_err(|error| MixedPortfolioReplayError::Instance {
970 instance_id: instance_id.clone(),
971 reason: error.to_string(),
972 })?;
973 configured_series.push(MultiTimeframeSeries::new(specs).map_err(|error| {
974 MixedPortfolioReplayError::Instance {
975 instance_id: instance_id.clone(),
976 reason: error.to_string(),
977 }
978 })?);
979 instance_ids.push(instance_id);
980 profiles.push(instance.entry_profiles);
981 configured_feed_from.push(instance.feed_from);
982 configured_analysis.push(instance.analysis);
983 adapters.push(instance.adapter);
984 }
985
986 let mut direct_strategies = Vec::with_capacity(direct.len());
987 let mut direct_analysis = Vec::with_capacity(direct.len());
988 let mut direct_series = Vec::with_capacity(direct.len());
989 let mut direct_feed_from = Vec::with_capacity(direct.len());
990 for instance in direct {
991 let instance_id = instance.instance_id;
992 if !identities.insert(instance_id.clone()) {
993 return Err(MixedPortfolioReplayError::DuplicateInstanceIdentity { instance_id });
994 }
995 crate::strategy::replay::validate_series_specs(
996 instance.strategy.requirements(),
997 &instance.series,
998 )
999 .map_err(|error| MixedPortfolioReplayError::Instance {
1000 instance_id: instance_id.clone(),
1001 reason: error.to_string(),
1002 })?;
1003 direct_series.push(MultiTimeframeSeries::new(instance.series).map_err(|error| {
1004 MixedPortfolioReplayError::Instance {
1005 instance_id: instance_id.clone(),
1006 reason: error.to_string(),
1007 }
1008 })?);
1009 instance_ids.push(instance_id);
1010 profiles.push(instance.entry_profiles);
1011 direct_feed_from.push(instance.feed_from);
1012 direct_analysis.push(instance.analysis);
1013 direct_strategies.push(instance.strategy);
1014 }
1015 self.instance_profiles = profiles;
1016
1017 let configured_count = adapters.len();
1018 let research_limits = self.strategy_research_limits;
1019 let configured_drivers = adapters
1020 .iter_mut()
1021 .zip(configured_series)
1022 .zip(configured_analysis)
1023 .map(|((adapter, series), analysis)| {
1024 ConfiguredStrategyReplayDriver::new(
1025 adapter,
1026 series,
1027 analysis,
1028 retention,
1029 research_limits,
1030 )
1031 })
1032 .collect::<Vec<_>>();
1033 let direct_drivers = direct_strategies
1034 .iter_mut()
1035 .zip(direct_series)
1036 .zip(direct_analysis)
1037 .map(|((strategy, series), analysis)| {
1038 StrategyReplayDriver::new(
1039 strategy.as_mut(),
1040 series,
1041 analysis,
1042 retention,
1043 research_limits,
1044 )
1045 })
1046 .collect::<Vec<_>>();
1047 let configured_declared = configured_drivers
1048 .iter()
1049 .map(|driver| declared_series(&driver.requirements))
1050 .collect();
1051 let direct_declared = direct_drivers
1052 .iter()
1053 .map(|driver| declared_series(&driver.requirements))
1054 .collect();
1055 let mut hook = MixedPortfolioReplayHook {
1056 configured: configured_drivers,
1057 direct: direct_drivers,
1058 configured_declared,
1059 direct_declared,
1060 configured_feed_from,
1061 direct_feed_from,
1062 state: PortfolioReplayState::new(instance_ids.clone(), supervisor),
1063 failed: None,
1064 seen_input: (false, false),
1065 mixed_input: None,
1066 };
1067 let replay = match self.run_raw_signals_future_batches(
1068 feed,
1069 primary_eod,
1070 Vec::new(),
1071 None,
1072 future,
1073 None,
1074 0,
1075 0,
1076 &mut is_cancelled,
1077 &mut on_progress,
1078 &mut hook,
1079 ) {
1080 Ok(replay) => replay,
1081 Err(FutureBatchReplayError::Feed(error)) => {
1082 return Err(MixedPortfolioReplayError::Feed(error));
1083 }
1084 Err(FutureBatchReplayError::Cancelled) => {
1085 return Err(MixedPortfolioReplayError::Cancelled);
1086 }
1087 Err(FutureBatchReplayError::Dynamic) => {
1088 if let Some(timestamp) = hook.mixed_input {
1089 return Err(MixedPortfolioReplayError::MixedPrimaryInput { timestamp });
1090 }
1091 let failure = hook
1092 .failed
1093 .take()
1094 .unwrap_or(MixedDriverFailure::Configured(0));
1095 let (instance, reason) = match failure {
1096 MixedDriverFailure::Configured(index) => {
1097 let reason = hook
1098 .configured
1099 .swap_remove(index)
1100 .finish()
1101 .err()
1102 .map_or_else(
1103 || "configured instance failed".into(),
1104 |error| format!("{error:?}"),
1105 );
1106 (index, reason)
1107 }
1108 MixedDriverFailure::Direct(index) => {
1109 let reason = hook.direct.swap_remove(index).finish().err().map_or_else(
1110 || "direct instance failed".into(),
1111 |error| mixed_direct_error(&error),
1112 );
1113 (configured_count + index, reason)
1114 }
1115 };
1116 return Err(MixedPortfolioReplayError::Instance {
1117 instance_id: instance_ids[instance].clone(),
1118 reason,
1119 });
1120 }
1121 };
1122 for (index, driver) in hook.configured.into_iter().enumerate() {
1123 driver
1124 .finish()
1125 .map_err(|error| MixedPortfolioReplayError::Instance {
1126 instance_id: instance_ids[index].clone(),
1127 reason: format!("{error:?}"),
1128 })?;
1129 }
1130 for (index, driver) in hook.direct.into_iter().enumerate() {
1131 driver
1132 .finish()
1133 .map_err(|error| MixedPortfolioReplayError::Instance {
1134 instance_id: instance_ids[configured_count + index].clone(),
1135 reason: mixed_direct_error(&error),
1136 })?;
1137 }
1138 Ok(MixedPortfolioBacktestResult {
1139 replay,
1140 supervisor: hook.state.into_output(),
1141 })
1142 }
1143
1144 #[allow(clippy::too_many_arguments)]
1148 pub fn run_portfolio_future_streaming_controlled<F, C, P>(
1149 mut self,
1150 feed: &mut F,
1151 primary_eod: Option<NaiveDateTime>,
1152 instances: Vec<ConfiguredInstance>,
1153 supervisor: Option<PortfolioSupervisor>,
1154 retention: StrategyRetentionLimits,
1155 mut is_cancelled: C,
1156 mut on_progress: P,
1157 ) -> Result<PortfolioBacktestResult, PortfolioReplayError<F::Error>>
1158 where
1159 F: FallibleBatchFeed,
1160 C: FnMut() -> bool,
1161 P: FnMut(ReplayProgress),
1162 {
1163 if instances.is_empty() {
1164 return Err(PortfolioReplayError::NoInstances);
1165 }
1166 if instances.len() > MAX_PORTFOLIO_INSTANCES {
1167 return Err(PortfolioReplayError::TooManyInstances(instances.len()));
1168 }
1169 let mut identities = BTreeSet::new();
1171 for instance in &instances {
1172 if !identities.insert(instance.instance_id().to_owned()) {
1173 return Err(PortfolioReplayError::DuplicateInstanceIdentity {
1174 instance_id: instance.instance_id().to_owned(),
1175 });
1176 }
1177 }
1178 if self.entry_profiles.is_some() {
1179 return Err(PortfolioReplayError::Input(
1180 StrategyReplayInputError::ManagementProfile(
1181 "portfolio instances carry their own entry profiles, so the runner must not set run-level routes".into(),
1182 ),
1183 ));
1184 }
1185 if self.config.run_tags.contains_key(INSTANCE_POSITION_TAG) {
1186 return Err(PortfolioReplayError::Input(
1187 StrategyReplayInputError::FutureQuote(format!(
1188 "run tag '{INSTANCE_POSITION_TAG}' is owned by the portfolio replay, which labels every position with its instance"
1189 )),
1190 ));
1191 }
1192 if let Some(supervisor) = supervisor.as_ref()
1193 && supervisor.caps_group_risk()
1194 && !self.config.sizing.as_ref().is_some_and(is_monetary_sizing)
1195 {
1196 return Err(PortfolioReplayError::Supervisor(
1197 "a group risk cap needs a monetary sizing policy, because a fixed-lot entry's risk is unknown until it fills".into(),
1198 ));
1199 }
1200 let future = self.future_config.clone().unwrap_or_default();
1201 self.future_config = Some(future.clone());
1202 validate_replay_config(&self.config, Some(&future), &[])
1203 .map_err(StrategyReplayInputError::FutureQuote)?;
1204
1205 let mut adapters = Vec::with_capacity(instances.len());
1206 let mut analyses = Vec::with_capacity(instances.len());
1207 let mut instance_ids = Vec::with_capacity(instances.len());
1208 let mut strategy_ids = Vec::with_capacity(instances.len());
1209 let mut profiles = Vec::with_capacity(instances.len());
1210 let mut feed_from = Vec::with_capacity(instances.len());
1211 let mut series = Vec::with_capacity(instances.len());
1212 for instance in instances {
1213 let instance_id = instance.instance_id().to_owned();
1214 let fail = |source: StrategyReplayInputError| PortfolioReplayError::Instance {
1215 instance_id: instance_id.clone(),
1216 source: StrategyReplayError::Input(source),
1217 };
1218 instance
1219 .adapter
1220 .preflight_entry_profiles(&instance.entry_profiles)
1221 .map_err(|error| fail(error.into()))?;
1222 instance
1223 .adapter
1224 .preflight(retention, self.strategy_research_limits)
1225 .map_err(|error| fail(error.into()))?;
1226 if !instance
1228 .adapter
1229 .configured_requirements()
1230 .entries
1231 .is_empty()
1232 {
1233 let symbol = instance.adapter.configured_strategy().primary_symbol();
1234 let problem = match self.config.sizing.as_ref() {
1235 None => Some(
1236 "a strategy that emits Entry actions requires a sizing policy".to_owned(),
1237 ),
1238 Some(_)
1239 if !self.config.symbol_specs.contains_key(symbol)
1240 && explicit_instrument_spec(&self.config, symbol).is_none() =>
1241 {
1242 Some(format!("missing instrument or symbol spec for {symbol}"))
1243 }
1244 Some(policy)
1245 if is_monetary_sizing(policy)
1246 && !future.currency_plan.as_ref().is_some_and(|plan| {
1247 plan.route_for_primary_symbol(symbol).is_some()
1248 }) =>
1249 {
1250 Some(format!(
1251 "monetary sizing needs a currency plan route for primary symbol {symbol}"
1252 ))
1253 }
1254 Some(_) => None,
1255 };
1256 if let Some(problem) = problem {
1257 return Err(fail(StrategyReplayInputError::FutureQuote(problem)));
1258 }
1259 }
1260 let specs = instance.adapter.series_specs().cloned().collect::<Vec<_>>();
1261 crate::strategy::replay::validate_series_specs(instance.adapter.requirements(), &specs)
1262 .map_err(fail)?;
1263 let instance_series = MultiTimeframeSeries::new(specs).map_err(|error| {
1264 PortfolioReplayError::Instance {
1265 instance_id: instance_id.clone(),
1266 source: StrategyReplayError::Series(error),
1267 }
1268 })?;
1269 strategy_ids.push(instance.strategy_id().to_owned());
1270 instance_ids.push(instance_id);
1271 profiles.push(instance.entry_profiles);
1272 feed_from.push(instance.feed_from);
1273 analyses.push(instance.analysis);
1274 series.push(instance_series);
1275 adapters.push(instance.adapter);
1276 }
1277 self.instance_profiles = profiles.clone();
1278
1279 let research_limits = self.strategy_research_limits;
1280 let descriptors = adapters
1281 .iter()
1282 .map(|adapter| adapter.descriptor().clone())
1283 .collect::<Vec<_>>();
1284 let drivers = adapters
1285 .iter_mut()
1286 .zip(series)
1287 .zip(analyses)
1288 .map(|((adapter, series), analysis)| {
1289 ConfiguredStrategyReplayDriver::new(
1290 adapter,
1291 series,
1292 analysis,
1293 retention,
1294 research_limits,
1295 )
1296 })
1297 .collect::<Vec<_>>();
1298 let declared = drivers
1299 .iter()
1300 .map(|driver| {
1301 let mut declared = BTreeMap::<String, BTreeSet<u64>>::new();
1302 for requirement in driver.requirements.series() {
1303 declared
1304 .entry(requirement.symbol().to_owned())
1305 .or_default()
1306 .insert(requirement.timeframe().duration_seconds());
1307 }
1308 declared
1309 })
1310 .collect();
1311 let mut hook = PortfolioReplayHook {
1312 drivers,
1313 declared,
1314 feed_from,
1315 state: PortfolioReplayState::new(instance_ids.clone(), supervisor),
1316 failed: None,
1317 seen_input: (false, false),
1318 mixed_input: None,
1319 };
1320 let replay = match self.run_raw_signals_future_batches(
1321 feed,
1322 primary_eod,
1323 Vec::new(),
1324 None,
1325 future,
1326 None,
1327 0,
1328 0,
1329 &mut is_cancelled,
1330 &mut on_progress,
1331 &mut hook,
1332 ) {
1333 Ok(replay) => replay,
1334 Err(FutureBatchReplayError::Feed(error)) => {
1335 return Err(PortfolioReplayError::Feed(error));
1336 }
1337 Err(FutureBatchReplayError::Cancelled) => return Err(PortfolioReplayError::Cancelled),
1338 Err(FutureBatchReplayError::Dynamic) => {
1339 if let Some(timestamp) = hook.mixed_input {
1340 return Err(PortfolioReplayError::MixedPrimaryInput { timestamp });
1341 }
1342 let index = hook.failed.unwrap_or(0);
1343 let driver = hook.drivers.swap_remove(index);
1344 let error = driver
1345 .finish()
1346 .expect_err("dynamic failure stores its cause");
1347 return Err(PortfolioReplayError::Instance {
1348 instance_id: instance_ids[index].clone(),
1349 source: map_strategy_driver_error(error),
1350 });
1351 }
1352 };
1353 let PortfolioReplayHook { drivers, state, .. } = hook;
1354 let mut outputs = Vec::with_capacity(drivers.len());
1355 for (index, driver) in drivers.into_iter().enumerate() {
1356 let (decisions, research) =
1357 driver
1358 .finish()
1359 .map_err(|error| PortfolioReplayError::Instance {
1360 instance_id: instance_ids[index].clone(),
1361 source: map_strategy_driver_error(error),
1362 })?;
1363 outputs.push(PortfolioInstanceOutput {
1364 strategy_id: strategy_ids[index].clone(),
1365 instance_id: instance_ids[index].clone(),
1366 descriptor: descriptors[index].clone(),
1367 decisions,
1368 research,
1369 });
1370 }
1371 Ok(PortfolioBacktestResult {
1372 replay,
1373 instances: outputs,
1374 supervisor: state.into_output(),
1375 })
1376 }
1377}
1378
1379fn mixed_direct_error<E: std::fmt::Display>(error: &StrategyDriverError<E>) -> String {
1380 match error {
1381 StrategyDriverError::Series(error) => error.to_string(),
1382 StrategyDriverError::SeriesView(error) => error.to_string(),
1383 StrategyDriverError::Analysis(error) => error.to_string(),
1384 StrategyDriverError::Strategy(error) => error.to_string(),
1385 StrategyDriverError::Runtime(error) => error.to_string(),
1386 StrategyDriverError::WarmupSignals { timestamp } => {
1387 format!("strategy emitted signals during warmup at {timestamp}")
1388 }
1389 StrategyDriverError::InvalidGeneratedSignal {
1390 signal_index,
1391 reason,
1392 } => format!("generated signal {signal_index} is invalid: {reason}"),
1393 StrategyDriverError::TickExecutionRequired { symbol, timestamp } => {
1394 format!("{symbol} requires tick execution at {timestamp}")
1395 }
1396 }
1397}
1398
1399fn declared_series(
1400 requirements: &crate::strategy::StrategyRequirements,
1401) -> BTreeMap<String, BTreeSet<u64>> {
1402 let mut declared = BTreeMap::<String, BTreeSet<u64>>::new();
1403 for requirement in requirements.series() {
1404 declared
1405 .entry(requirement.symbol().to_owned())
1406 .or_default()
1407 .insert(requirement.timeframe().duration_seconds());
1408 }
1409 declared
1410}
1411
1412fn requested_account_risk(
1414 sizing: Option<&SizingPolicy>,
1415 balance: f64,
1416 risk_multiplier: f64,
1417) -> Option<f64> {
1418 let risk = match sizing? {
1419 SizingPolicy::FixedRiskAmount { amount } => amount * risk_multiplier,
1420 SizingPolicy::BalanceRiskPercent { percent } => balance * percent / 100.0 * risk_multiplier,
1421 SizingPolicy::FixedLot { .. } => return None,
1422 };
1423 (risk.is_finite() && risk > 0.0).then_some(risk)
1424}