1use std::collections::{BTreeMap, BTreeSet, HashMap, VecDeque};
18use std::convert::Infallible;
19
20use chrono::{Duration, NaiveDateTime};
21use qs_core::sizing::{
22 compute_instrument_native_loss_per_lot, compute_instrument_size_for_spec_with_prices,
23};
24use qs_core::types::{
25 Action, CloseReason, Effect, ExecutionFill, ExecutionModel, FillModel, FutureEffect, OrderType,
26 PositionStatus, PreparedPendingFill, PriceQuote, Side, SlippageModel, position_size_tolerance,
27};
28use qs_core::{ExecutionPricer, FutureApplyError, TradeEngine};
29use qs_instruments::{
30 Decimal, DecimalGrid, EconomicsModelId, InstrumentSpec, ListingStatus, PositiveDecimal,
31 QuantityUnit,
32};
33use serde::Serialize;
34
35use crate::artifacts::{
36 EntryProfileResolutionAudit, EntryProfileSelectionSource, EntryResolutionStage,
37 ExecutionMetadata, FUTURE_ARTIFACT_FORMAT_VERSION, FutureBacktestArtifacts,
38 InstrumentSizingArtifact, MarketEntrySizingAudit, MarketEntrySizingBasis, PendingOrderSnapshot,
39 ReplayInstrumentManifest,
40};
41use crate::currency::{ConversionQuoteBook, RunCurrencyPlan};
42use crate::data_feed::{
43 BarExecutionPrices, DataFeed, FallibleBatchFeed, FeedEvent, MarketEvent, TimestampBatch,
44};
45use crate::economic_support::{LEGACY_ECONOMIC_GUARD_ID, resolve_legacy_economics};
46use crate::evaluation::EvaluationOptions;
47use crate::executor::BacktestExecutor;
48use crate::future_executor::{FutureExecutor, FutureExecutorError};
49use crate::ledger::{ActionDisposition, ActionDispositionStatus, LifecycleLedger};
50use crate::mtm::{MtmCurveCollector, MtmOutputPolicy, MtmOutputSummary};
51use crate::portfolio::{EquityPoint, PortfolioRecorder};
52use crate::profile::{
53 EntryProfileRoutingError, EntryResolutionContext, ManagementProfile, PreparedEntryProfiles,
54 PriceGridSource, RawSignal, ResolvedEntry, allocate_target_steps, resolve_signal,
55 resolve_unprofiled_entry,
56};
57use crate::report::BacktestResult;
58use crate::sizing::{SizingPolicy, compute_native_loss_per_lot, compute_size};
59use crate::strategy::configured::BoundaryPositionFacts;
60use crate::strategy::{
61 AnalysisBoundary, AnalysisPipeline, BacktestConfiguredStrategyAdapter, BarSeriesSpec,
62 ConfiguredStrategyAdapterError, HistoricalStrategy, MultiTimeframeSeries, Strategy,
63 StrategyBacktestResult, StrategyContext, StrategyDecisionRecorder, StrategyEvent,
64 StrategyFeedback, StrategyFeedbackEvent, StrategyJournalRecorder, StrategyReplayError,
65 StrategyReplayInputError, StrategyResearchLimits, StrategyResearchOutput,
66 StrategyRetentionLimits,
67};
68
69#[derive(Debug, Clone, Serialize)]
72pub struct FutureQuoteConfig {
73 pub signal_latency_ms: i64,
75 pub slippage_pips: f64,
77 pub stale_quote_after_ms: Option<i64>,
79 pub pnl_epsilon: f64,
81 pub currency_plan: Option<RunCurrencyPlan>,
83 pub conversion_stale_after_ms: i64,
85 pub mtm_output: MtmOutputPolicy,
87 pub market_entry_sizing_basis: MarketEntrySizingBasis,
89}
90
91impl Default for FutureQuoteConfig {
92 fn default() -> Self {
93 Self {
94 signal_latency_ms: 0,
95 slippage_pips: 0.0,
96 stale_quote_after_ms: None,
97 pnl_epsilon: 1.0e-9,
98 currency_plan: None,
99 conversion_stale_after_ms: 300_000,
100 mtm_output: MtmOutputPolicy::default(),
101 market_entry_sizing_basis: MarketEntrySizingBasis::default(),
102 }
103 }
104}
105
106#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
108pub struct ReplayProgress {
109 pub processed_events: usize,
110 pub total_events: usize,
111 pub processed_signals: usize,
112 pub total_signals: usize,
113}
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
117#[error("backtest replay cancelled")]
118pub struct ReplayCancelled;
119
120#[derive(Debug, thiserror::Error)]
122pub enum StreamingReplayError<E> {
123 #[error("market-data stream failed: {0}")]
124 Feed(E),
125 #[error(transparent)]
126 Cancelled(#[from] ReplayCancelled),
127}
128
129const REPLAY_PROGRESS_INTERVAL: usize = 256;
130
131mod portfolio_replay;
132
133const MAX_BAR_LEG_STEPS: usize = 100_000;
135
136#[derive(Debug, Clone, Copy)]
138enum QuoteSide {
139 Bid,
140 Ask,
141 Mid,
142}
143
144fn should_report_progress(processed: usize, total: usize) -> bool {
145 processed == total || processed.is_multiple_of(REPLAY_PROGRESS_INTERVAL)
146}
147
148#[derive(Debug, Clone)]
149pub(crate) struct ScheduledSignal {
150 sequence: u64,
151 signal_ts: NaiveDateTime,
152 effective_ts: NaiveDateTime,
153 signal: RawSignal,
154 action_id: Option<String>,
155 explicit_action_id: bool,
157 requires_later_quote: bool,
158 instance: Option<usize>,
160}
161
162impl ScheduledSignal {
163 pub(crate) fn new(
164 sequence: u64,
165 signal_ts: NaiveDateTime,
166 effective_ts: NaiveDateTime,
167 signal: RawSignal,
168 requires_later_quote: bool,
169 ) -> Self {
170 Self {
171 sequence,
172 signal_ts,
173 effective_ts,
174 signal,
175 action_id: None,
176 explicit_action_id: false,
177 requires_later_quote,
178 instance: None,
179 }
180 }
181
182 #[allow(dead_code)]
183 pub(crate) fn with_action_id(mut self, action_id: impl Into<String>) -> Self {
184 self.action_id = Some(action_id.into());
185 self.explicit_action_id = true;
186 self
187 }
188
189 fn with_action_base(mut self, base: impl Into<String>) -> Self {
191 self.action_id = Some(base.into());
192 self.explicit_action_id = false;
193 self
194 }
195
196 fn resolved_action_id(&self) -> String {
197 self.action_id
198 .clone()
199 .unwrap_or_else(|| format!("signal:{:08}", self.sequence))
200 }
201}
202
203#[derive(Debug, Clone)]
204struct QueuedAction {
205 action_id: String,
206 action_kind: String,
207 action: Action,
208 execution: Option<ExecutionFill>,
209 symbol: String,
210 signal_ts: NaiveDateTime,
211 effective_ts: NaiveDateTime,
212 entry_signal: Option<RawSignal>,
213 entry_profile: Option<ManagementProfile>,
214 entry_profile_selection_source: Option<EntryProfileSelectionSource>,
215 selected_profile_name: Option<String>,
216 market_entry_sizing_audit: Option<MarketEntrySizingAudit>,
217 entry_profile_resolution_audit: Option<EntryProfileResolutionAudit>,
218 requires_later_quote: bool,
219}
220
221#[derive(Debug, Clone)]
222struct SelectedEntryProfile {
223 profile: Option<ManagementProfile>,
224 source: EntryProfileSelectionSource,
225 entry_class: Option<String>,
226 profile_name: Option<String>,
227}
228
229struct FinalizedEntry {
230 action: Action,
231 requested_account_risk: Option<f64>,
232 native_loss_per_lot: Option<f64>,
233 account_loss_per_lot: Option<f64>,
234 final_lot: f64,
235 level_resolution: qs_core::EntryLevelResolution,
236 target_resolution: qs_core::TargetResolution,
237 configured_weights: Vec<f64>,
238 allocated_target_steps: Vec<u64>,
239 remainder_steps: u64,
240}
241
242#[derive(Debug, thiserror::Error)]
243enum FutureTransactionError {
244 #[error(transparent)]
245 Core(#[from] FutureApplyError),
246 #[error(transparent)]
247 Accounting(#[from] FutureExecutorError),
248}
249
250#[derive(Debug)]
251enum StrategyDriverError<E> {
252 Series(crate::strategy::SeriesError),
253 SeriesView(crate::strategy::SeriesViewError),
254 Analysis(crate::strategy::AnalysisError),
255 Strategy(E),
256 Runtime(crate::strategy::StrategyRuntimeError),
257 WarmupSignals {
258 timestamp: NaiveDateTime,
259 },
260 InvalidGeneratedSignal {
261 signal_index: usize,
262 reason: String,
263 },
264 TickExecutionRequired {
265 symbol: String,
266 timestamp: NaiveDateTime,
267 },
268}
269
270enum FutureBatchReplayError<E> {
271 Feed(E),
272 Cancelled,
273 Dynamic,
274}
275
276trait FutureReplayHook {
277 fn is_active(&self) -> bool;
278 fn output_ready(&self) -> bool;
279 fn preflight_primary_events(&mut self, events: &[FeedEvent]) -> bool;
280 fn reject_generated_configuration(&mut self, instance: Option<usize>, reason: String);
281 fn observes_position_economics(&self) -> bool;
283 fn bar_execution_timeframes(&self) -> BTreeMap<String, u64> {
285 BTreeMap::new()
286 }
287 fn reads_completed_bars_only(&self) -> bool {
289 false
290 }
291 fn retains_post_bar_boundary(&self) -> bool {
293 false
294 }
295 #[allow(clippy::too_many_arguments)]
296 fn on_pre_bar_boundary(
297 &mut self,
298 batch: &TimestampBatch,
299 engine: &TradeEngine,
300 lifecycle: &LifecycleLedger,
301 positions: &BoundaryPositionFacts<'_>,
302 pending_effects: &mut Vec<FutureEffect>,
303 pending_events: &mut Vec<StrategyFeedbackEvent>,
304 ) -> Option<Vec<ScheduledSignal>> {
305 self.on_boundary(
306 batch,
307 engine,
308 lifecycle,
309 positions,
310 pending_effects,
311 pending_events,
312 )
313 }
314 fn portfolio_state(&mut self) -> Option<&mut portfolio_replay::PortfolioReplayState> {
316 None
317 }
318 #[allow(clippy::too_many_arguments)]
319 fn on_boundary(
320 &mut self,
321 batch: &TimestampBatch,
322 engine: &TradeEngine,
323 lifecycle: &LifecycleLedger,
324 positions: &BoundaryPositionFacts<'_>,
325 pending_effects: &mut Vec<FutureEffect>,
326 pending_events: &mut Vec<StrategyFeedbackEvent>,
327 ) -> Option<Vec<ScheduledSignal>>;
328 fn on_final_committed(
329 &mut self,
330 pending_effects: &mut Vec<FutureEffect>,
331 pending_events: &mut Vec<StrategyFeedbackEvent>,
332 ) -> bool;
333}
334
335fn shortest_series_per_symbol(
337 requirements: &crate::strategy::StrategyRequirements,
338) -> BTreeMap<String, u64> {
339 let mut shortest = BTreeMap::<String, u64>::new();
340 for series in requirements.series() {
341 let seconds = series.timeframe().duration_seconds();
342 shortest
343 .entry(series.symbol().to_owned())
344 .and_modify(|current| *current = (*current).min(seconds))
345 .or_insert(seconds);
346 }
347 shortest
348}
349
350struct StaticReplayHook;
351
352impl FutureReplayHook for StaticReplayHook {
353 fn observes_position_economics(&self) -> bool {
354 false
355 }
356
357 fn is_active(&self) -> bool {
358 false
359 }
360
361 fn output_ready(&self) -> bool {
362 true
363 }
364
365 fn preflight_primary_events(&mut self, _events: &[FeedEvent]) -> bool {
366 true
367 }
368
369 fn reject_generated_configuration(&mut self, _instance: Option<usize>, _reason: String) {
370 unreachable!("static replay does not generate strategy signals");
371 }
372
373 fn on_boundary(
374 &mut self,
375 _batch: &TimestampBatch,
376 _engine: &TradeEngine,
377 _lifecycle: &LifecycleLedger,
378 _positions: &BoundaryPositionFacts<'_>,
379 pending_effects: &mut Vec<FutureEffect>,
380 pending_events: &mut Vec<StrategyFeedbackEvent>,
381 ) -> Option<Vec<ScheduledSignal>> {
382 pending_effects.clear();
383 pending_events.clear();
384 Some(Vec::new())
385 }
386
387 fn on_final_committed(
388 &mut self,
389 pending_effects: &mut Vec<FutureEffect>,
390 pending_events: &mut Vec<StrategyFeedbackEvent>,
391 ) -> bool {
392 pending_effects.clear();
393 pending_events.clear();
394 true
395 }
396}
397
398struct StrategyReplayDriver<'a, S: HistoricalStrategy + ?Sized> {
399 strategy: &'a mut S,
400 requirements: crate::strategy::StrategyRequirements,
401 series: MultiTimeframeSeries,
402 analysis: AnalysisPipeline,
403 limits: StrategyRetentionLimits,
404 decisions: StrategyDecisionRecorder,
405 journal: StrategyJournalRecorder,
406 next_decision_sequence: u64,
407 next_signal_sequence: u64,
408 delivered_dispositions: usize,
409 warmup_complete: bool,
410 failure: Option<StrategyDriverError<S::Error>>,
411}
412
413impl<'a, S: HistoricalStrategy + ?Sized> StrategyReplayDriver<'a, S> {
414 fn new(
415 strategy: &'a mut S,
416 series: MultiTimeframeSeries,
417 analysis: AnalysisPipeline,
418 limits: StrategyRetentionLimits,
419 research_limits: StrategyResearchLimits,
420 ) -> Self {
421 Self {
422 requirements: strategy.requirements().clone(),
423 strategy,
424 series,
425 analysis,
426 limits,
427 decisions: StrategyDecisionRecorder::new(limits),
428 journal: StrategyJournalRecorder::new(research_limits),
429 next_decision_sequence: 0,
430 next_signal_sequence: 0,
431 delivered_dispositions: 0,
432 warmup_complete: false,
433 failure: None,
434 }
435 }
436
437 fn finish(
438 self,
439 ) -> Result<
440 (
441 crate::strategy::StrategyDecisionOutput,
442 StrategyResearchOutput,
443 ),
444 StrategyDriverError<S::Error>,
445 > {
446 match self.failure {
447 Some(error) => Err(error),
448 None => Ok((
449 self.decisions.finish(),
450 StrategyResearchOutput {
451 journal: self.journal.finish(),
452 research_annotations: self.analysis.into_research_annotations(),
453 },
454 )),
455 }
456 }
457
458 fn fail(&mut self, error: StrategyDriverError<S::Error>) -> Option<Vec<ScheduledSignal>> {
459 self.failure = Some(error);
460 None
461 }
462}
463
464impl<S: HistoricalStrategy + ?Sized> FutureReplayHook for StrategyReplayDriver<'_, S> {
465 fn observes_position_economics(&self) -> bool {
466 false
467 }
468
469 fn bar_execution_timeframes(&self) -> BTreeMap<String, u64> {
470 shortest_series_per_symbol(&self.requirements)
471 }
472
473 fn is_active(&self) -> bool {
474 true
475 }
476
477 fn output_ready(&self) -> bool {
478 self.warmup_complete
479 }
480
481 fn preflight_primary_events(&mut self, events: &[FeedEvent]) -> bool {
482 if self.requirements.needs_tick_execution()
483 && let Some(event) = events
484 .iter()
485 .find(|event| matches!(event.event, MarketEvent::Bar { .. }))
486 {
487 self.failure = Some(StrategyDriverError::TickExecutionRequired {
488 symbol: event.event.symbol().to_owned(),
489 timestamp: event.event.ts(),
490 });
491 return false;
492 }
493 true
494 }
495
496 fn reject_generated_configuration(&mut self, _instance: Option<usize>, reason: String) {
497 self.failure = Some(StrategyDriverError::InvalidGeneratedSignal {
498 signal_index: 0,
499 reason,
500 });
501 }
502
503 fn on_boundary(
504 &mut self,
505 batch: &TimestampBatch,
506 engine: &TradeEngine,
507 lifecycle: &LifecycleLedger,
508 _positions: &BoundaryPositionFacts<'_>,
509 pending_effects: &mut Vec<FutureEffect>,
510 pending_events: &mut Vec<StrategyFeedbackEvent>,
511 ) -> Option<Vec<ScheduledSignal>> {
512 let closed_bars = match self.series.on_batch(batch) {
513 Ok(bars) => bars,
514 Err(error) => return self.fail(StrategyDriverError::Series(error)),
515 };
516 let boundary = AnalysisBoundary::new(batch.ts, &closed_bars, &self.series);
517 let observations = match self.analysis.on_boundary(boundary) {
518 Ok(output) => output.observations().to_vec(),
519 Err(error) => return self.fail(StrategyDriverError::Analysis(error)),
520 };
521 self.warmup_complete = match self.series.warmup_complete(&self.requirements) {
522 Ok(complete) => complete,
523 Err(error) => return self.fail(StrategyDriverError::SeriesView(error)),
524 };
525 let disposition_end = lifecycle.len();
526 let feedback = StrategyFeedback::with_events(
527 pending_effects,
528 &lifecycle.as_slice()[self.delivered_dispositions..disposition_end],
529 pending_events,
530 );
531 let event = StrategyEvent::new(&batch.events, &closed_bars, &observations, feedback);
532 let context = StrategyContext::new(
533 batch.ts,
534 &self.series,
535 self.analysis.observations(),
536 engine,
537 self.warmup_complete,
538 );
539 let output = match self.strategy.on_event(event, context) {
540 Ok(output) => output,
541 Err(error) => return self.fail(StrategyDriverError::Strategy(error)),
542 };
543 pending_effects.clear();
544 pending_events.clear();
545 self.delivered_dispositions = disposition_end;
546
547 let (decision, journal) = output.into_parts();
548 if let Err(error) = self.journal.push_callback(batch.ts, journal) {
549 return self.fail(StrategyDriverError::Runtime(
550 crate::strategy::StrategyRuntimeError::Journal(error),
551 ));
552 }
553 let Some(draft) = decision else {
554 return Some(Vec::new());
555 };
556 let record = match draft.into_record(self.next_decision_sequence, batch.ts, self.limits) {
557 Ok(record) => record,
558 Err(error) => return self.fail(StrategyDriverError::Runtime(error)),
559 };
560 if !self.warmup_complete && !record.emitted_signals().is_empty() {
561 return self.fail(StrategyDriverError::WarmupSignals {
562 timestamp: batch.ts,
563 });
564 }
565 if let Some((signal_index, error)) =
566 record
567 .emitted_signals()
568 .iter()
569 .enumerate()
570 .find_map(|(index, signal)| {
571 qs_core::validation::validate_raw_signal(signal)
572 .err()
573 .map(|error| (index, error))
574 })
575 {
576 return self.fail(StrategyDriverError::InvalidGeneratedSignal {
577 signal_index,
578 reason: error.to_string(),
579 });
580 }
581 let effective_ts = match self.requirements.effective_timestamp(batch.ts) {
582 Ok(timestamp) => timestamp,
583 Err(error) => {
584 return self.fail(StrategyDriverError::Runtime(
585 crate::strategy::StrategyRuntimeError::Domain(error),
586 ));
587 }
588 };
589 let signals = match self.decisions.push(record) {
590 Ok(signals) => signals,
591 Err(error) => {
592 return self.fail(StrategyDriverError::Runtime(
593 crate::strategy::StrategyRuntimeError::Domain(error),
594 ));
595 }
596 };
597 self.next_decision_sequence = match self.next_decision_sequence.checked_add(1) {
598 Some(sequence) => sequence,
599 None => {
600 return self.fail(StrategyDriverError::Runtime(
601 crate::strategy::StrategyRuntimeError::Domain(
602 crate::strategy::StrategyDomainError::OmittedCounterOverflow,
603 ),
604 ));
605 }
606 };
607 let mut scheduled = Vec::with_capacity(signals.len());
608 for signal in signals {
609 let sequence = self.next_signal_sequence;
610 self.next_signal_sequence = match self.next_signal_sequence.checked_add(1) {
611 Some(sequence) => sequence,
612 None => {
613 return self.fail(StrategyDriverError::Runtime(
614 crate::strategy::StrategyRuntimeError::Domain(
615 crate::strategy::StrategyDomainError::OmittedCounterOverflow,
616 ),
617 ));
618 }
619 };
620 scheduled.push(ScheduledSignal::new(
621 sequence,
622 batch.ts,
623 effective_ts,
624 signal,
625 true,
626 ));
627 }
628 Some(scheduled)
629 }
630
631 fn on_final_committed(
632 &mut self,
633 pending_effects: &mut Vec<FutureEffect>,
634 pending_events: &mut Vec<StrategyFeedbackEvent>,
635 ) -> bool {
636 pending_effects.clear();
637 pending_events.clear();
638 true
639 }
640}
641
642struct ConfiguredStrategyReplayDriver<'a> {
643 adapter: &'a mut BacktestConfiguredStrategyAdapter,
644 requirements: crate::strategy::StrategyRequirements,
645 series: MultiTimeframeSeries,
646 analysis: AnalysisPipeline,
647 limits: StrategyRetentionLimits,
648 research_limits: StrategyResearchLimits,
649 decisions: StrategyDecisionRecorder,
650 journal: StrategyJournalRecorder,
651 next_decision_sequence: u64,
652 next_signal_sequence: u64,
653 warmup_complete: bool,
654 failure: Option<StrategyDriverError<ConfiguredStrategyAdapterError>>,
655}
656
657impl<'a> ConfiguredStrategyReplayDriver<'a> {
658 fn new(
659 adapter: &'a mut BacktestConfiguredStrategyAdapter,
660 series: MultiTimeframeSeries,
661 analysis: AnalysisPipeline,
662 limits: StrategyRetentionLimits,
663 research_limits: StrategyResearchLimits,
664 ) -> Self {
665 Self {
666 requirements: adapter.requirements().clone(),
667 adapter,
668 series,
669 analysis,
670 limits,
671 research_limits,
672 decisions: StrategyDecisionRecorder::new(limits),
673 journal: StrategyJournalRecorder::new(research_limits),
674 next_decision_sequence: 0,
675 next_signal_sequence: 0,
676 warmup_complete: false,
677 failure: None,
678 }
679 }
680
681 fn finish(
682 self,
683 ) -> Result<
684 (
685 crate::strategy::StrategyDecisionOutput,
686 StrategyResearchOutput,
687 ),
688 StrategyDriverError<ConfiguredStrategyAdapterError>,
689 > {
690 match self.failure {
691 Some(error) => Err(error),
692 None => Ok((
693 self.decisions.finish(),
694 StrategyResearchOutput {
695 journal: self.journal.finish(),
696 research_annotations: self.analysis.into_research_annotations(),
697 },
698 )),
699 }
700 }
701
702 fn fail(
703 &mut self,
704 error: StrategyDriverError<ConfiguredStrategyAdapterError>,
705 ) -> Option<Vec<ScheduledSignal>> {
706 self.failure = Some(error);
707 None
708 }
709}
710
711impl FutureReplayHook for ConfiguredStrategyReplayDriver<'_> {
712 fn observes_position_economics(&self) -> bool {
713 true
714 }
715
716 fn bar_execution_timeframes(&self) -> BTreeMap<String, u64> {
717 shortest_series_per_symbol(&self.requirements)
718 }
719
720 fn reads_completed_bars_only(&self) -> bool {
721 true
722 }
723
724 fn is_active(&self) -> bool {
725 true
726 }
727
728 fn output_ready(&self) -> bool {
729 self.warmup_complete
730 }
731
732 fn preflight_primary_events(&mut self, events: &[FeedEvent]) -> bool {
734 if !self
735 .adapter
736 .configured_requirements()
737 .count_required_sources
738 .is_empty()
739 && events.iter().any(|event| {
740 matches!(
741 event.event,
742 MarketEvent::Bar {
743 tick_count: None,
744 ..
745 }
746 )
747 })
748 {
749 self.reject_generated_configuration(None, "configured strategy requires tick count but a primary stored bar has unknown count".into());
750 false
751 } else {
752 true
753 }
754 }
755
756 fn reject_generated_configuration(&mut self, _instance: Option<usize>, reason: String) {
757 self.failure = Some(StrategyDriverError::InvalidGeneratedSignal {
758 signal_index: 0,
759 reason,
760 });
761 }
762
763 fn on_boundary(
764 &mut self,
765 batch: &TimestampBatch,
766 engine: &TradeEngine,
767 _lifecycle: &LifecycleLedger,
768 positions: &BoundaryPositionFacts<'_>,
769 pending_effects: &mut Vec<FutureEffect>,
770 pending_events: &mut Vec<StrategyFeedbackEvent>,
771 ) -> Option<Vec<ScheduledSignal>> {
772 let closed_bars = match self.series.on_batch(batch) {
773 Ok(bars) => bars,
774 Err(error) => return self.fail(StrategyDriverError::Series(error)),
775 };
776 let boundary = AnalysisBoundary::new(batch.ts, &closed_bars, &self.series);
777 let observations = match self.analysis.on_boundary(boundary) {
778 Ok(output) => output.observations().to_vec(),
779 Err(error) => return self.fail(StrategyDriverError::Analysis(error)),
780 };
781 self.warmup_complete = match self.series.warmup_complete(&self.requirements) {
782 Ok(complete) => complete,
783 Err(error) => return self.fail(StrategyDriverError::SeriesView(error)),
784 };
785 let output = match self.adapter.evaluate_boundary(
786 batch.ts,
787 self.warmup_complete,
788 &closed_bars,
789 &observations,
790 &self.series,
791 self.analysis.observations(),
792 engine,
793 positions,
794 pending_events,
795 self.limits,
796 self.research_limits,
797 ) {
798 Ok(output) => output,
799 Err(error) => return self.fail(StrategyDriverError::Strategy(error)),
800 };
801 pending_effects.clear();
802 pending_events.clear();
803
804 if let Err(error) = self.journal.push_callback(batch.ts, output.journal) {
805 return self.fail(StrategyDriverError::Runtime(
806 crate::strategy::StrategyRuntimeError::Journal(error),
807 ));
808 }
809 if let Some(decision) = output.decision {
810 let record =
811 match decision.into_record(self.next_decision_sequence, batch.ts, self.limits) {
812 Ok(record) => record,
813 Err(error) => return self.fail(StrategyDriverError::Runtime(error)),
814 };
815 if let Err(error) = self.decisions.push(record) {
816 return self.fail(StrategyDriverError::Runtime(
817 crate::strategy::StrategyRuntimeError::Domain(error),
818 ));
819 }
820 self.next_decision_sequence = match self.next_decision_sequence.checked_add(1) {
821 Some(sequence) => sequence,
822 None => {
823 return self.fail(StrategyDriverError::Runtime(
824 crate::strategy::StrategyRuntimeError::Domain(
825 crate::strategy::StrategyDomainError::OmittedCounterOverflow,
826 ),
827 ));
828 }
829 };
830 }
831
832 if !self.warmup_complete && !output.commands.is_empty() {
833 return self.fail(StrategyDriverError::WarmupSignals {
834 timestamp: batch.ts,
835 });
836 }
837 let effective_ts = match self.requirements.effective_timestamp(batch.ts) {
838 Ok(timestamp) => timestamp,
839 Err(error) => {
840 return self.fail(StrategyDriverError::Runtime(
841 crate::strategy::StrategyRuntimeError::Domain(error),
842 ));
843 }
844 };
845 let mut scheduled = Vec::with_capacity(output.commands.len());
846 for (signal_index, command) in output.commands.into_iter().enumerate() {
847 if let Err(error) = qs_core::validation::validate_raw_signal(&command.signal) {
848 return self.fail(StrategyDriverError::InvalidGeneratedSignal {
849 signal_index,
850 reason: error.to_string(),
851 });
852 }
853 let sequence = self.next_signal_sequence;
854 self.next_signal_sequence = match self.next_signal_sequence.checked_add(1) {
855 Some(sequence) => sequence,
856 None => {
857 return self.fail(StrategyDriverError::Runtime(
858 crate::strategy::StrategyRuntimeError::Domain(
859 crate::strategy::StrategyDomainError::OmittedCounterOverflow,
860 ),
861 ));
862 }
863 };
864 scheduled.push(
865 ScheduledSignal::new(sequence, batch.ts, effective_ts, command.signal, true)
866 .with_action_id(command.command_id),
867 );
868 }
869 Some(scheduled)
870 }
871
872 fn on_final_committed(
873 &mut self,
874 pending_effects: &mut Vec<FutureEffect>,
875 pending_events: &mut Vec<StrategyFeedbackEvent>,
876 ) -> bool {
877 pending_effects.clear();
878 match self.adapter.finalize_feedback(pending_events) {
879 Ok(()) => {
880 pending_events.clear();
881 true
882 }
883 Err(error) => {
884 self.failure = Some(StrategyDriverError::Strategy(error));
885 false
886 }
887 }
888 }
889}
890
891#[derive(Debug, Clone, Copy, PartialEq, Eq)]
893pub enum EquityObservationKind {
894 PreSettlement,
895 PostOutput,
896 ConversionRevaluation,
897 QuiescentTermination,
898 EndOfData,
899}
900
901impl EquityObservationKind {
902 pub const fn as_str(self) -> &'static str {
903 match self {
904 Self::PreSettlement => "pre_settlement",
905 Self::PostOutput => "post_output",
906 Self::ConversionRevaluation => "conversion_revaluation",
907 Self::QuiescentTermination => "quiescent_termination",
908 Self::EndOfData => "end_of_data",
909 }
910 }
911}
912
913struct BufferedFutureFeed {
914 events: VecDeque<FeedEvent>,
915 total_events: usize,
916}
917
918impl BufferedFutureFeed {
919 fn new(events: Vec<FeedEvent>) -> Self {
920 let total_events = events.len();
921 Self {
922 events: VecDeque::from(events),
923 total_events,
924 }
925 }
926}
927
928impl DataFeed for BufferedFutureFeed {
929 fn next_event(&mut self) -> Option<MarketEvent> {
930 self.events.pop_front().map(|event| event.event)
931 }
932
933 fn peek(&self) -> Option<&MarketEvent> {
934 self.events.front().map(|event| &event.event)
935 }
936
937 fn next_batch(&mut self) -> Option<TimestampBatch> {
938 let ts = self.events.front()?.event.ts();
939 let mut events = Vec::new();
940 while self
941 .events
942 .front()
943 .is_some_and(|event| event.event.ts() == ts)
944 {
945 events.push(self.events.pop_front().expect("front checked"));
946 }
947 Some(TimestampBatch { ts, events })
948 }
949
950 fn total_events(&self) -> Option<usize> {
951 Some(self.total_events)
952 }
953}
954
955struct DataFeedBatchAdapter<'a, F> {
956 feed: &'a mut F,
957}
958
959impl<F: DataFeed> FallibleBatchFeed for DataFeedBatchAdapter<'_, F> {
960 type Error = Infallible;
961
962 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
963 Ok(self.feed.next_batch())
964 }
965}
966
967#[derive(Debug, Clone, Serialize)]
969pub struct BacktestConfig {
970 pub initial_balance: f64,
972 pub close_on_finish: bool,
975 pub fill_model: FillModel,
980 pub contract_sizes: HashMap<String, f64>,
989 pub sizing: Option<SizingPolicy>,
992 pub symbol_specs: HashMap<String, qs_symbols::SymbolSpec>,
995 pub instrument_manifest: Option<ReplayInstrumentManifest>,
997 pub bar_spread_fallback: HashMap<String, f64>,
1001 pub run_tags: BTreeMap<String, String>,
1005 pub costs: HashMap<String, qs_core::InstrumentCosts>,
1009}
1010
1011impl Default for BacktestConfig {
1012 fn default() -> Self {
1013 Self {
1014 initial_balance: 10_000.0,
1015 close_on_finish: true,
1016 fill_model: FillModel::default(),
1017 contract_sizes: HashMap::new(),
1018 sizing: None,
1019 symbol_specs: HashMap::new(),
1020 instrument_manifest: None,
1021 bar_spread_fallback: HashMap::new(),
1022 run_tags: BTreeMap::new(),
1023 costs: HashMap::new(),
1024 }
1025 }
1026}
1027
1028pub struct BacktestRunner {
1030 engine: TradeEngine,
1031 executor: BacktestExecutor,
1032 config: BacktestConfig,
1033 future_config: Option<FutureQuoteConfig>,
1034 evaluation_options: EvaluationOptions,
1035 strategy_research_limits: StrategyResearchLimits,
1036 instrument_sizing: Vec<InstrumentSizingArtifact>,
1037 market_entry_sizing: Vec<MarketEntrySizingAudit>,
1038 entry_profile_resolutions: Vec<EntryProfileResolutionAudit>,
1039 committed_feedback: Vec<FutureEffect>,
1040 committed_feedback_events: Vec<StrategyFeedbackEvent>,
1041 entry_profiles: Option<PreparedEntryProfiles>,
1042 instance_profiles: Vec<PreparedEntryProfiles>,
1044}
1045
1046impl BacktestRunner {
1047 pub fn new(config: BacktestConfig) -> Self {
1049 let executor =
1050 BacktestExecutor::new(config.initial_balance, effective_contract_sizes(&config));
1051 Self {
1052 engine: TradeEngine::with_fill_model(config.fill_model),
1053 executor,
1054 config,
1055 future_config: None,
1056 evaluation_options: EvaluationOptions::default(),
1057 strategy_research_limits: StrategyResearchLimits::default(),
1058 instrument_sizing: Vec::new(),
1059 market_entry_sizing: Vec::new(),
1060 entry_profile_resolutions: Vec::new(),
1061 committed_feedback: Vec::new(),
1062 committed_feedback_events: Vec::new(),
1063 entry_profiles: None,
1064 instance_profiles: Vec::new(),
1065 }
1066 }
1067
1068 pub fn new_future(config: BacktestConfig, future_config: FutureQuoteConfig) -> Self {
1070 let executor =
1071 BacktestExecutor::new(config.initial_balance, effective_contract_sizes(&config));
1072 let engine = TradeEngine::with_fill_model_and_deterministic_ids(config.fill_model);
1073 Self {
1074 engine,
1075 executor,
1076 config,
1077 future_config: Some(future_config),
1078 evaluation_options: EvaluationOptions::default(),
1079 strategy_research_limits: StrategyResearchLimits::default(),
1080 instrument_sizing: Vec::new(),
1081 market_entry_sizing: Vec::new(),
1082 entry_profile_resolutions: Vec::new(),
1083 committed_feedback: Vec::new(),
1084 committed_feedback_events: Vec::new(),
1085 entry_profiles: None,
1086 instance_profiles: Vec::new(),
1087 }
1088 }
1089
1090 pub fn with_defaults() -> Self {
1092 Self::new(BacktestConfig::default())
1093 }
1094
1095 pub fn with_entry_profiles(mut self, profiles: PreparedEntryProfiles) -> Self {
1097 self.entry_profiles = Some(profiles);
1098 self
1099 }
1100
1101 pub fn with_evaluation_options(mut self, options: EvaluationOptions) -> Self {
1104 self.evaluation_options = options;
1105 self
1106 }
1107
1108 pub fn with_strategy_research_limits(mut self, limits: StrategyResearchLimits) -> Self {
1110 self.strategy_research_limits = limits;
1111 self
1112 }
1113
1114 pub fn engine(&self) -> &TradeEngine {
1116 &self.engine
1117 }
1118
1119 pub fn executor(&self) -> &BacktestExecutor {
1121 &self.executor
1122 }
1123
1124 pub fn run_strategy<F: DataFeed, S: Strategy>(
1138 mut self,
1139 feed: &mut F,
1140 strategy: &mut S,
1141 ) -> BacktestResult {
1142 if validate_replay_config(&self.config, None, &[]).is_err() {
1143 return rejected_legacy_result(&self.config);
1144 }
1145 let mut last_quote_ts = BTreeMap::new();
1146 while let Some(event) = feed.next_event() {
1147 let quote = event.to_quote();
1148 if !accept_legacy_quote("e, &mut last_quote_ts) {
1149 continue;
1150 }
1151
1152 let price_effects = self.engine.on_price("e);
1154 self.executor
1155 .process_effects(&price_effects, &self.engine, "e);
1156
1157 let actions = strategy.on_event(&event);
1159 self.apply_actions(actions, "e);
1160 }
1161
1162 let final_actions = strategy.on_finished();
1164 if !final_actions.is_empty() {
1165 if let Some(last_quote) = self.last_available_quote() {
1169 self.apply_actions(final_actions, &last_quote);
1170 let effects = self.engine.on_price(&last_quote);
1172 self.executor
1173 .process_effects(&effects, &self.engine, &last_quote);
1174 }
1175 }
1176
1177 self.close_remaining_if_configured();
1179
1180 BacktestResult::from_trade_log(self.config.initial_balance, self.executor.trade_log)
1181 }
1182
1183 fn apply_actions(&mut self, actions: Vec<Action>, quote: &PriceQuote) {
1187 for action in actions {
1188 self.apply_single_action(action, quote.ts, quote);
1189 }
1190 }
1191
1192 fn apply_single_action(
1199 &mut self,
1200 action: Action,
1201 ts: chrono::NaiveDateTime,
1202 quote: &PriceQuote,
1203 ) {
1204 match self.engine.apply_action(action, ts) {
1205 Ok(effects) => {
1206 self.executor.process_effects(&effects, &self.engine, quote);
1207 }
1208 Err(_) => {
1209 }
1213 }
1214 }
1215
1216 fn last_available_quote(&self) -> Option<PriceQuote> {
1218 for pos in self.engine.open_positions() {
1221 if let Some(q) = self.engine.last_quote(&pos.data.symbol) {
1222 return Some(q.clone());
1223 }
1224 }
1225 for pos in self.engine.closed_positions() {
1227 if let Some(q) = self.engine.last_quote(&pos.data.symbol) {
1228 return Some(q.clone());
1229 }
1230 }
1231 None
1232 }
1233
1234 pub fn run_raw_signals<F: DataFeed>(
1243 self,
1244 feed: &mut F,
1245 raw_signals: Vec<RawSignal>,
1246 profile: Option<&ManagementProfile>,
1247 ) -> BacktestResult {
1248 self.run_raw_signals_controlled(feed, raw_signals, profile, || false, |_| {})
1249 .expect("non-cancellable replay cannot be cancelled")
1250 }
1251
1252 pub fn run_raw_signals_controlled<F, C, P>(
1258 mut self,
1259 feed: &mut F,
1260 raw_signals: Vec<RawSignal>,
1261 profile: Option<&ManagementProfile>,
1262 mut is_cancelled: C,
1263 mut on_progress: P,
1264 ) -> std::result::Result<BacktestResult, ReplayCancelled>
1265 where
1266 F: DataFeed,
1267 C: FnMut() -> bool,
1268 P: FnMut(ReplayProgress),
1269 {
1270 if let Some(future_config) = self.future_config.clone() {
1271 return self.run_raw_signals_future_controlled(
1272 feed,
1273 raw_signals,
1274 profile,
1275 future_config,
1276 &mut is_cancelled,
1277 &mut on_progress,
1278 );
1279 }
1280 if validate_replay_config(&self.config, None, &raw_signals).is_err()
1281 || profile.is_some_and(|profile| profile.validate().is_err())
1282 || self.validate_entry_profile_routes(&raw_signals).is_err()
1283 {
1284 return Ok(rejected_legacy_result(&self.config));
1285 }
1286
1287 let total_events = feed.total_events().unwrap_or(0);
1288 let total_signals = raw_signals.len();
1289 let mut processed_events = 0;
1290 let mut sig_idx = 0;
1291 let mut last_quote_ts = BTreeMap::new();
1292 on_progress(ReplayProgress {
1293 processed_events,
1294 total_events,
1295 processed_signals: sig_idx,
1296 total_signals,
1297 });
1298
1299 while let Some(event) = feed.next_event() {
1300 if is_cancelled() {
1301 return Err(ReplayCancelled);
1302 }
1303 let quote = event.to_quote();
1304 if !accept_legacy_quote("e, &mut last_quote_ts) {
1305 processed_events += 1;
1306 if should_report_progress(processed_events, total_events) {
1307 on_progress(ReplayProgress {
1308 processed_events,
1309 total_events,
1310 processed_signals: sig_idx,
1311 total_signals,
1312 });
1313 }
1314 continue;
1315 }
1316
1317 while sig_idx < raw_signals.len() && raw_signals[sig_idx].ts() <= event.ts() {
1319 if is_cancelled() {
1320 return Err(ReplayCancelled);
1321 }
1322 self.process_raw_signal(&raw_signals[sig_idx], profile, "e);
1323 sig_idx += 1;
1324 if should_report_progress(sig_idx, total_signals) {
1325 on_progress(ReplayProgress {
1326 processed_events,
1327 total_events,
1328 processed_signals: sig_idx,
1329 total_signals,
1330 });
1331 }
1332 }
1333
1334 let effects = self.engine.on_price("e);
1336 self.executor
1337 .process_effects(&effects, &self.engine, "e);
1338 processed_events += 1;
1339 if should_report_progress(processed_events, total_events) {
1340 on_progress(ReplayProgress {
1341 processed_events,
1342 total_events,
1343 processed_signals: sig_idx,
1344 total_signals,
1345 });
1346 }
1347 }
1348
1349 if sig_idx < raw_signals.len()
1351 && let Some(last_quote) = self.last_available_quote()
1352 {
1353 while sig_idx < raw_signals.len() {
1354 if is_cancelled() {
1355 return Err(ReplayCancelled);
1356 }
1357 self.process_raw_signal(&raw_signals[sig_idx], profile, &last_quote);
1358 sig_idx += 1;
1359 if should_report_progress(sig_idx, total_signals) {
1360 on_progress(ReplayProgress {
1361 processed_events,
1362 total_events,
1363 processed_signals: sig_idx,
1364 total_signals,
1365 });
1366 }
1367 }
1368 let effects = self.engine.on_price(&last_quote);
1370 self.executor
1371 .process_effects(&effects, &self.engine, &last_quote);
1372 }
1373
1374 if is_cancelled() {
1375 return Err(ReplayCancelled);
1376 }
1377
1378 self.close_remaining_if_configured();
1380 on_progress(ReplayProgress {
1381 processed_events,
1382 total_events,
1383 processed_signals: sig_idx,
1384 total_signals,
1385 });
1386
1387 Ok(BacktestResult::from_trade_log(
1388 self.config.initial_balance,
1389 self.executor.trade_log,
1390 ))
1391 }
1392
1393 pub fn run_raw_signals_future<F: DataFeed>(
1399 self,
1400 feed: &mut F,
1401 raw_signals: Vec<RawSignal>,
1402 profile: Option<&ManagementProfile>,
1403 ) -> BacktestResult {
1404 let future_config = self.future_config.clone().unwrap_or_default();
1405 self.run_raw_signals_future_with_config(feed, raw_signals, profile, future_config)
1406 }
1407
1408 #[allow(clippy::too_many_arguments)]
1412 pub fn run_raw_signals_future_streaming_controlled<F, C, P>(
1413 mut self,
1414 feed: &mut F,
1415 primary_eod: Option<NaiveDateTime>,
1416 raw_signals: Vec<RawSignal>,
1417 profile: Option<&ManagementProfile>,
1418 mut is_cancelled: C,
1419 mut on_progress: P,
1420 ) -> std::result::Result<BacktestResult, StreamingReplayError<F::Error>>
1421 where
1422 F: FallibleBatchFeed,
1423 C: FnMut() -> bool,
1424 P: FnMut(ReplayProgress),
1425 {
1426 let future = self.future_config.clone().unwrap_or_default();
1427 self.future_config = Some(future.clone());
1428 if let Err(error) = validate_replay_config(&self.config, Some(&future), &raw_signals) {
1429 return Ok(rejected_future_result(
1430 &self.config,
1431 &future,
1432 self.evaluation_options,
1433 error,
1434 ));
1435 }
1436 if let Some(profile) = profile
1437 && let Err(error) = profile.validate()
1438 {
1439 return Ok(rejected_future_result(
1440 &self.config,
1441 &future,
1442 self.evaluation_options,
1443 error.to_string(),
1444 ));
1445 }
1446 if let Err(error) = self.validate_entry_profile_routes(&raw_signals) {
1447 return Ok(rejected_future_result(
1448 &self.config,
1449 &future,
1450 self.evaluation_options,
1451 error,
1452 ));
1453 }
1454
1455 let mut hook = StaticReplayHook;
1456 match self.run_raw_signals_future_batches(
1457 feed,
1458 primary_eod,
1459 raw_signals,
1460 profile,
1461 future,
1462 None,
1463 0,
1464 0,
1465 &mut is_cancelled,
1466 &mut on_progress,
1467 &mut hook,
1468 ) {
1469 Ok(result) => Ok(result),
1470 Err(FutureBatchReplayError::Feed(error)) => Err(StreamingReplayError::Feed(error)),
1471 Err(FutureBatchReplayError::Cancelled) => {
1472 Err(StreamingReplayError::Cancelled(ReplayCancelled))
1473 }
1474 Err(FutureBatchReplayError::Dynamic) => {
1475 unreachable!("static replay has no dynamic hook")
1476 }
1477 }
1478 }
1479
1480 pub fn run_raw_signals_future_fallible<F: FallibleBatchFeed>(
1484 self,
1485 feed: &mut F,
1486 raw_signals: Vec<RawSignal>,
1487 profile: Option<&ManagementProfile>,
1488 ) -> Result<BacktestResult, F::Error> {
1489 let future_config = self.future_config.clone().unwrap_or_default();
1490 self.run_raw_signals_future_fallible_with_config(feed, raw_signals, profile, future_config)
1491 }
1492
1493 pub fn run_raw_signals_future_fallible_with_config<F: FallibleBatchFeed>(
1495 self,
1496 feed: &mut F,
1497 raw_signals: Vec<RawSignal>,
1498 profile: Option<&ManagementProfile>,
1499 future: FutureQuoteConfig,
1500 ) -> Result<BacktestResult, F::Error> {
1501 if let Err(error) = validate_replay_config(&self.config, Some(&future), &raw_signals) {
1502 return Ok(rejected_future_result(
1503 &self.config,
1504 &future,
1505 self.evaluation_options,
1506 error,
1507 ));
1508 }
1509 if let Some(profile) = profile
1510 && let Err(error) = profile.validate()
1511 {
1512 return Ok(rejected_future_result(
1513 &self.config,
1514 &future,
1515 self.evaluation_options,
1516 error.to_string(),
1517 ));
1518 }
1519 if let Err(error) = self.validate_entry_profile_routes(&raw_signals) {
1520 return Ok(rejected_future_result(
1521 &self.config,
1522 &future,
1523 self.evaluation_options,
1524 error,
1525 ));
1526 }
1527
1528 let mut events = Vec::new();
1529 while let Some(batch) = FallibleBatchFeed::next_batch(feed)? {
1530 events.extend(batch.events);
1531 }
1532 let mut buffered_feed = BufferedFutureFeed::new(events);
1533 Ok(self.run_raw_signals_future_with_config(
1534 &mut buffered_feed,
1535 raw_signals,
1536 profile,
1537 future,
1538 ))
1539 }
1540
1541 #[allow(clippy::too_many_arguments)]
1543 pub fn run_historical_strategy_future<F, S>(
1544 self,
1545 source_feed: &mut F,
1546 strategy: &mut S,
1547 series_specs: Vec<BarSeriesSpec>,
1548 analysis: AnalysisPipeline,
1549 retention: StrategyRetentionLimits,
1550 profile: Option<&ManagementProfile>,
1551 ) -> Result<StrategyBacktestResult, StrategyReplayError<Infallible, S::Error>>
1552 where
1553 F: DataFeed,
1554 S: HistoricalStrategy + ?Sized,
1555 {
1556 crate::strategy::replay::validate_series_specs(strategy.requirements(), &series_specs)?;
1557 MultiTimeframeSeries::new(series_specs.clone())?;
1558 let future = self.future_config.clone().unwrap_or_default();
1559 validate_replay_config(&self.config, Some(&future), &[])
1560 .map_err(StrategyReplayInputError::FutureQuote)?;
1561 if let Some(profile) = profile {
1562 profile
1563 .validate()
1564 .map_err(|error| StrategyReplayInputError::ManagementProfile(error.to_string()))?;
1565 }
1566 let mut ordered_events = Vec::new();
1567 let mut source_last_ts = BTreeMap::<String, NaiveDateTime>::new();
1568 while let Some(batch) = source_feed.next_batch() {
1569 for event in batch.events {
1570 let symbol = event.event.symbol().to_owned();
1571 let timestamp = event.event.ts();
1572 if source_last_ts
1573 .get(&symbol)
1574 .is_some_and(|previous| *previous > timestamp)
1575 {
1576 continue;
1577 }
1578 source_last_ts.insert(symbol, timestamp);
1579 ordered_events.push(event);
1580 }
1581 }
1582 ordered_events.sort_by_key(FeedEvent::ordering_key);
1583 let primary_eod = ordered_events
1584 .iter()
1585 .filter(|event| event.metadata.roles.primary)
1586 .filter_map(|event| event.event.to_valid_quote())
1587 .map(|quote| quote.ts)
1588 .max();
1589 let mut ordered_feed = crate::data_feed::VecFeed::from_feed_events(ordered_events);
1590 let mut feed = DataFeedBatchAdapter {
1591 feed: &mut ordered_feed,
1592 };
1593 self.run_historical_strategy_future_streaming(
1594 &mut feed,
1595 primary_eod,
1596 strategy,
1597 series_specs,
1598 analysis,
1599 retention,
1600 profile,
1601 )
1602 }
1603
1604 #[allow(clippy::too_many_arguments)]
1606 pub fn run_historical_strategy_future_streaming<F, S>(
1607 mut self,
1608 feed: &mut F,
1609 primary_eod: Option<NaiveDateTime>,
1610 strategy: &mut S,
1611 series_specs: Vec<BarSeriesSpec>,
1612 analysis: AnalysisPipeline,
1613 retention: StrategyRetentionLimits,
1614 profile: Option<&ManagementProfile>,
1615 ) -> Result<StrategyBacktestResult, StrategyReplayError<F::Error, S::Error>>
1616 where
1617 F: FallibleBatchFeed,
1618 S: HistoricalStrategy + ?Sized,
1619 {
1620 crate::strategy::replay::validate_series_specs(strategy.requirements(), &series_specs)?;
1621 let future = self.future_config.clone().unwrap_or_default();
1622 self.future_config = Some(future.clone());
1623 validate_replay_config(&self.config, Some(&future), &[])
1624 .map_err(StrategyReplayInputError::FutureQuote)?;
1625 if let Some(profile) = profile {
1626 profile
1627 .validate()
1628 .map_err(|error| StrategyReplayInputError::ManagementProfile(error.to_string()))?;
1629 }
1630 let descriptor = strategy.descriptor().clone();
1631 let series = MultiTimeframeSeries::new(series_specs)?;
1632 let mut hook = StrategyReplayDriver::new(
1633 strategy,
1634 series,
1635 analysis,
1636 retention,
1637 self.strategy_research_limits,
1638 );
1639 let mut is_cancelled = || false;
1640 let mut on_progress = |_| {};
1641 let replay = match self.run_raw_signals_future_batches(
1642 feed,
1643 primary_eod,
1644 Vec::new(),
1645 profile,
1646 future,
1647 None,
1648 0,
1649 0,
1650 &mut is_cancelled,
1651 &mut on_progress,
1652 &mut hook,
1653 ) {
1654 Ok(replay) => replay,
1655 Err(FutureBatchReplayError::Feed(error)) => {
1656 return Err(StrategyReplayError::Feed(error));
1657 }
1658 Err(FutureBatchReplayError::Cancelled) => {
1659 unreachable!("strategy replay is not cancellable")
1660 }
1661 Err(FutureBatchReplayError::Dynamic) => {
1662 let error = hook.finish().expect_err("dynamic failure stores its cause");
1663 return Err(map_strategy_driver_error(error));
1664 }
1665 };
1666 let (decisions, research) = hook.finish().map_err(map_strategy_driver_error)?;
1667 Ok(StrategyBacktestResult {
1668 replay,
1669 descriptor,
1670 decisions,
1671 research,
1672 })
1673 }
1674
1675 #[allow(clippy::too_many_arguments)]
1677 pub fn run_configured_strategy_future<F>(
1678 self,
1679 source_feed: &mut F,
1680 adapter: &mut BacktestConfiguredStrategyAdapter,
1681 analysis: AnalysisPipeline,
1682 retention: StrategyRetentionLimits,
1683 profile: Option<&ManagementProfile>,
1684 ) -> Result<
1685 StrategyBacktestResult,
1686 StrategyReplayError<Infallible, ConfiguredStrategyAdapterError>,
1687 >
1688 where
1689 F: DataFeed,
1690 {
1691 self.preflight_configured_entry_profiles(adapter, profile)?;
1692 adapter
1693 .preflight(retention, self.strategy_research_limits)
1694 .map_err(StrategyReplayInputError::ConfiguredAdapter)?;
1695 let series_specs = adapter.series_specs().cloned().collect::<Vec<_>>();
1696 crate::strategy::replay::validate_series_specs(adapter.requirements(), &series_specs)?;
1697 MultiTimeframeSeries::new(series_specs.clone())?;
1698 let future = self.future_config.clone().unwrap_or_default();
1699 validate_replay_config(&self.config, Some(&future), &[])
1700 .map_err(StrategyReplayInputError::FutureQuote)?;
1701
1702 let mut ordered_events = Vec::new();
1703 let mut source_last_ts = BTreeMap::<String, NaiveDateTime>::new();
1704 while let Some(batch) = source_feed.next_batch() {
1705 for event in batch.events {
1706 let symbol = event.event.symbol().to_owned();
1707 let timestamp = event.event.ts();
1708 if source_last_ts
1709 .get(&symbol)
1710 .is_some_and(|previous| *previous > timestamp)
1711 {
1712 continue;
1713 }
1714 source_last_ts.insert(symbol, timestamp);
1715 ordered_events.push(event);
1716 }
1717 }
1718 ordered_events.sort_by_key(FeedEvent::ordering_key);
1719 let primary_eod = ordered_events
1720 .iter()
1721 .filter(|event| event.metadata.roles.primary)
1722 .filter_map(|event| event.event.to_valid_quote())
1723 .map(|quote| quote.ts)
1724 .max();
1725 let mut ordered_feed = crate::data_feed::VecFeed::from_feed_events(ordered_events);
1726 let mut feed = DataFeedBatchAdapter {
1727 feed: &mut ordered_feed,
1728 };
1729 self.run_configured_strategy_future_streaming(
1730 &mut feed,
1731 primary_eod,
1732 adapter,
1733 analysis,
1734 retention,
1735 profile,
1736 )
1737 }
1738
1739 #[allow(clippy::too_many_arguments)]
1741 pub fn run_configured_strategy_future_streaming<F>(
1742 self,
1743 feed: &mut F,
1744 primary_eod: Option<NaiveDateTime>,
1745 adapter: &mut BacktestConfiguredStrategyAdapter,
1746 analysis: AnalysisPipeline,
1747 retention: StrategyRetentionLimits,
1748 profile: Option<&ManagementProfile>,
1749 ) -> Result<StrategyBacktestResult, StrategyReplayError<F::Error, ConfiguredStrategyAdapterError>>
1750 where
1751 F: FallibleBatchFeed,
1752 {
1753 self.run_configured_strategy_future_streaming_controlled(
1754 feed,
1755 primary_eod,
1756 adapter,
1757 analysis,
1758 retention,
1759 profile,
1760 || false,
1761 |_| {},
1762 )
1763 }
1764
1765 #[allow(clippy::too_many_arguments)]
1767 pub fn run_configured_strategy_future_streaming_controlled<F, C, P>(
1768 mut self,
1769 feed: &mut F,
1770 primary_eod: Option<NaiveDateTime>,
1771 adapter: &mut BacktestConfiguredStrategyAdapter,
1772 analysis: AnalysisPipeline,
1773 retention: StrategyRetentionLimits,
1774 profile: Option<&ManagementProfile>,
1775 mut is_cancelled: C,
1776 mut on_progress: P,
1777 ) -> Result<StrategyBacktestResult, StrategyReplayError<F::Error, ConfiguredStrategyAdapterError>>
1778 where
1779 F: FallibleBatchFeed,
1780 C: FnMut() -> bool,
1781 P: FnMut(ReplayProgress),
1782 {
1783 self.preflight_configured_entry_profiles(adapter, profile)?;
1784 adapter
1785 .preflight(retention, self.strategy_research_limits)
1786 .map_err(StrategyReplayInputError::ConfiguredAdapter)?;
1787 let series_specs = adapter.series_specs().cloned().collect::<Vec<_>>();
1788 crate::strategy::replay::validate_series_specs(adapter.requirements(), &series_specs)?;
1789 let future = self.future_config.clone().unwrap_or_default();
1790 self.future_config = Some(future.clone());
1791 validate_replay_config(&self.config, Some(&future), &[])
1792 .map_err(StrategyReplayInputError::FutureQuote)?;
1793 let descriptor = adapter.descriptor().clone();
1794 let series = MultiTimeframeSeries::new(series_specs)?;
1795 let mut hook = ConfiguredStrategyReplayDriver::new(
1796 adapter,
1797 series,
1798 analysis,
1799 retention,
1800 self.strategy_research_limits,
1801 );
1802 let replay = match self.run_raw_signals_future_batches(
1803 feed,
1804 primary_eod,
1805 Vec::new(),
1806 profile,
1807 future,
1808 None,
1809 0,
1810 0,
1811 &mut is_cancelled,
1812 &mut on_progress,
1813 &mut hook,
1814 ) {
1815 Ok(replay) => replay,
1816 Err(FutureBatchReplayError::Feed(error)) => {
1817 return Err(StrategyReplayError::Feed(error));
1818 }
1819 Err(FutureBatchReplayError::Cancelled) => return Err(StrategyReplayError::Cancelled),
1820 Err(FutureBatchReplayError::Dynamic) => {
1821 let error = hook.finish().expect_err("dynamic failure stores its cause");
1822 return Err(map_strategy_driver_error(error));
1823 }
1824 };
1825 let (decisions, research) = hook.finish().map_err(map_strategy_driver_error)?;
1826 Ok(StrategyBacktestResult {
1827 replay,
1828 descriptor,
1829 decisions,
1830 research,
1831 })
1832 }
1833
1834 fn preflight_configured_entry_profiles(
1836 &self,
1837 adapter: &BacktestConfiguredStrategyAdapter,
1838 profile: Option<&ManagementProfile>,
1839 ) -> Result<(), StrategyReplayInputError> {
1840 if let Some(profile) = profile {
1841 profile
1842 .validate()
1843 .map_err(|error| StrategyReplayInputError::ManagementProfile(error.to_string()))?;
1844 }
1845 let run_default;
1846 let profiles = match self.entry_profiles.as_ref() {
1847 Some(profiles) => profiles,
1848 None => {
1849 run_default = PreparedEntryProfiles::default_only(profile.cloned());
1850 &run_default
1851 }
1852 };
1853 adapter.preflight_entry_profiles(profiles)?;
1854 Ok(())
1855 }
1856
1857 fn validate_entry_profile_routes(&self, signals: &[RawSignal]) -> Result<(), String> {
1858 if let Some(profiles) = self.entry_profiles.as_ref() {
1859 return profiles
1860 .validate_signals(signals)
1861 .map_err(|error| error.to_string());
1862 }
1863 if let Some(entry_class) = signals.iter().find_map(|signal| match signal {
1864 RawSignal::Entry {
1865 entry_class: Some(entry_class),
1866 ..
1867 } => Some(entry_class),
1868 _ => None,
1869 }) {
1870 return Err(
1871 EntryProfileRoutingError::UnknownEntryClass(entry_class.clone()).to_string(),
1872 );
1873 }
1874 Ok(())
1875 }
1876
1877 fn select_entry_profile(
1878 &self,
1879 signal: &RawSignal,
1880 fallback: Option<&ManagementProfile>,
1881 instance: Option<usize>,
1882 ) -> Result<SelectedEntryProfile, EntryProfileRoutingError> {
1883 let entry_class = match signal {
1884 RawSignal::Entry { entry_class, .. } => entry_class.clone(),
1885 _ => None,
1886 };
1887 let instance_profiles = instance.and_then(|instance| self.instance_profiles.get(instance));
1888 if let Some(profiles) = instance_profiles.or(self.entry_profiles.as_ref()) {
1889 let profile = profiles.select(signal)?.cloned();
1890 let source = if entry_class.is_some() {
1891 EntryProfileSelectionSource::Mapped
1892 } else if profile.is_some() {
1893 EntryProfileSelectionSource::RunDefault
1894 } else {
1895 EntryProfileSelectionSource::Unprofiled
1896 };
1897 let profile_name = profile.as_ref().map(|profile| profile.name.clone());
1898 return Ok(SelectedEntryProfile {
1899 profile,
1900 source,
1901 entry_class,
1902 profile_name,
1903 });
1904 }
1905 if let Some(entry_class) = entry_class {
1906 return Err(EntryProfileRoutingError::UnknownEntryClass(entry_class));
1907 }
1908 let profile = fallback.cloned();
1909 let source = if profile.is_some() {
1910 EntryProfileSelectionSource::RunDefault
1911 } else {
1912 EntryProfileSelectionSource::Unprofiled
1913 };
1914 let profile_name = profile.as_ref().map(|profile| profile.name.clone());
1915 Ok(SelectedEntryProfile {
1916 profile,
1917 source,
1918 entry_class: None,
1919 profile_name,
1920 })
1921 }
1922
1923 fn process_raw_signal(
1926 &mut self,
1927 signal: &RawSignal,
1928 profile: Option<&ManagementProfile>,
1929 quote: &PriceQuote,
1930 ) {
1931 let ts = signal.ts();
1932
1933 if signal.is_entry() {
1934 let mut signal = signal.clone();
1935 if let RawSignal::Entry {
1936 side,
1937 order_type: OrderType::Market,
1938 price,
1939 ..
1940 } = &mut signal
1941 && price.is_none()
1942 {
1943 *price = Some(match side {
1944 Side::Buy => quote.ask,
1945 Side::Sell => quote.bid,
1946 });
1947 }
1948 let selected_profile = match self.select_entry_profile(&signal, profile, None) {
1949 Ok(profile) => profile,
1950 Err(_) => return,
1951 };
1952 let resolved = match selected_profile.profile.as_ref() {
1953 Some(profile) => self.resolve_profiled_entry(profile, &signal),
1954 None => resolve_unprofiled_entry(&signal),
1955 };
1956 if let Ok(Some(resolved)) = resolved
1957 && let Ok(action) =
1958 self.finalize_resolved_entry(resolved, self.executor.balance, ts, None, None)
1959 {
1960 self.apply_single_action(action.action, ts, quote);
1961 }
1962 } else {
1963 let actions = resolve_signal(signal, &self.engine);
1964 for action in actions {
1965 self.apply_single_action(action, ts, quote);
1966 }
1967 }
1968 }
1969
1970 fn resolve_profiled_entry(
1971 &self,
1972 profile: &ManagementProfile,
1973 signal: &RawSignal,
1974 ) -> Result<Option<ResolvedEntry>, qs_core::ProfileApplicationError> {
1975 let symbol = match signal {
1976 RawSignal::Entry { symbol, .. } => symbol,
1977 _ => return profile.apply_entry_signal(signal),
1978 };
1979 match self.entry_resolution_context(symbol) {
1980 Ok(context) => profile.apply_entry_signal_with_context(signal, context),
1981 Err(_) => profile.apply_entry_signal(signal),
1982 }
1983 }
1984
1985 fn entry_resolution_context(&self, symbol: &str) -> Result<EntryResolutionContext, String> {
1986 if let Some(spec) = explicit_instrument_spec(&self.config, symbol) {
1987 return Ok(EntryResolutionContext {
1988 price_grid: spec.price.grid,
1989 price_grid_source: PriceGridSource::InstrumentPriceGrid,
1990 });
1991 }
1992 let legacy = self
1993 .config
1994 .symbol_specs
1995 .get(symbol)
1996 .ok_or_else(|| format!("missing price grid for {symbol}"))?;
1997 let scale = u8::try_from(legacy.digits)
1998 .map_err(|_| format!("price scale is too large for {symbol}"))?;
1999 let step = Decimal::new(1, scale).map_err(|error| error.to_string())?;
2000 let step = PositiveDecimal::new(step).map_err(|error| error.to_string())?;
2001 Ok(EntryResolutionContext {
2002 price_grid: DecimalGrid::new(Decimal::ZERO, step),
2003 price_grid_source: PriceGridSource::LegacyDigitsFallback,
2004 })
2005 }
2006
2007 fn finalize_resolved_entry(
2008 &mut self,
2009 mut resolved: ResolvedEntry,
2010 balance_before: f64,
2011 operation_ts: NaiveDateTime,
2012 conversion_quotes: Option<&ConversionQuoteBook>,
2013 sizing_reference_price: Option<f64>,
2014 ) -> Result<FinalizedEntry, String> {
2015 let policy = self
2016 .config
2017 .sizing
2018 .as_ref()
2019 .ok_or_else(|| "raw entry requires a sizing policy".to_owned())?;
2020 let entry_price = resolved.price.ok_or_else(|| {
2021 if resolved.order_type == OrderType::Market {
2022 "market entry requires an execution price".to_owned()
2023 } else {
2024 "pending entry requires a requested price".to_owned()
2025 }
2026 })?;
2027 let sizing_reference_price = sizing_reference_price.unwrap_or(entry_price);
2028 let explicit_spec = explicit_instrument_spec(&self.config, &resolved.symbol);
2029 let legacy_spec = self.config.symbol_specs.get(&resolved.symbol);
2030 if explicit_spec.is_none() && legacy_spec.is_none() {
2031 return Err(format!(
2032 "missing instrument or symbol spec for {}",
2033 resolved.symbol
2034 ));
2035 }
2036
2037 let (account_loss_per_lot, native_to_account_rate) = if is_monetary_sizing(policy) {
2038 let stop = resolved
2039 .stoploss
2040 .ok_or_else(|| "monetary sizing requires a protective stop".to_owned())?;
2041 let native_loss = match explicit_spec {
2042 Some(spec) => compute_instrument_native_loss_per_lot(
2043 resolved.side,
2044 sizing_reference_price,
2045 stop,
2046 u16::from(spec.price.display_scale),
2047 &spec.economics,
2048 )
2049 .map_err(|error| error.to_string())?,
2050 None => compute_native_loss_per_lot(
2051 resolved.side,
2052 sizing_reference_price,
2053 stop,
2054 legacy_spec.expect("legacy spec presence checked"),
2055 )
2056 .map_err(|error| error.to_string())?,
2057 };
2058 let plan = self
2059 .future_config
2060 .as_ref()
2061 .and_then(|config| config.currency_plan.as_ref())
2062 .ok_or_else(|| "monetary sizing requires a FutureQuote currency plan".to_owned())?;
2063 let route = plan
2064 .route_for_primary_symbol(&resolved.symbol)
2065 .ok_or_else(|| {
2066 format!(
2067 "currency plan has no frozen route for primary symbol {}",
2068 resolved.symbol
2069 )
2070 })?;
2071 let converted = conversion_quotes
2072 .ok_or_else(|| "monetary sizing requires conversion quotes".to_owned())?
2073 .convert_route(-native_loss, operation_ts, route)
2074 .map_err(|error| error.to_string())?;
2075 let account_loss = -converted.output_amount;
2076 (Some(account_loss), Some(account_loss / native_loss))
2077 } else {
2078 (None, None)
2079 };
2080
2081 let sizing = match explicit_spec {
2082 Some(spec) => compute_instrument_size_for_spec_with_prices(
2083 policy,
2084 resolved.risk_multiplier,
2085 balance_before,
2086 resolved.side,
2087 sizing_reference_price,
2088 entry_price,
2089 resolved.stoploss,
2090 spec,
2091 native_to_account_rate,
2092 )
2093 .map_err(|error| error.to_string())?,
2094 None => compute_size(
2095 policy,
2096 resolved.risk_multiplier,
2097 balance_before,
2098 resolved.side,
2099 sizing_reference_price,
2100 resolved.stoploss,
2101 legacy_spec.expect("legacy spec presence checked"),
2102 account_loss_per_lot,
2103 )
2104 .map_err(|error| error.to_string())?,
2105 };
2106 if let Some(quantity) = sizing.quantity_adjustment {
2107 self.instrument_sizing.push(InstrumentSizingArtifact {
2108 symbol: resolved.symbol.clone(),
2109 operation_ts,
2110 quantity,
2111 final_notional: sizing.final_notional.clone(),
2112 });
2113 }
2114 let configured_weights = resolved.target_resolution.weights.clone();
2115 let target_resolution = resolved.target_resolution.clone();
2116 let level_resolution = resolved.level_resolution.clone();
2117 let target_steps = allocate_target_steps(
2118 sizing.final_lot_steps,
2119 &configured_weights,
2120 resolved.target_resolution.remainder,
2121 )
2122 .map_err(|error| error.to_string())?;
2123 if target_steps.len() != resolved.targets.len() {
2124 return Err("target allocation does not match resolved targets".to_owned());
2125 }
2126 for (target, steps) in resolved.targets.iter_mut().zip(&target_steps) {
2127 target.close_ratio = *steps as f64 / sizing.final_lot_steps as f64;
2128 }
2129 let allocated_steps: u64 = target_steps.iter().sum();
2130 let remainder_steps = sizing.final_lot_steps.saturating_sub(allocated_steps);
2131
2132 Ok(FinalizedEntry {
2133 action: resolved.into_action(sizing.final_lot),
2134 requested_account_risk: sizing.requested_account_risk,
2135 native_loss_per_lot: sizing.native_loss_per_lot,
2136 account_loss_per_lot: sizing.account_loss_per_lot,
2137 final_lot: sizing.final_lot,
2138 level_resolution,
2139 target_resolution,
2140 configured_weights,
2141 allocated_target_steps: target_steps,
2142 remainder_steps,
2143 })
2144 }
2145
2146 fn close_remaining_if_configured(&mut self) {
2152 if !self.config.close_on_finish {
2153 return;
2154 }
2155
2156 let open_ids: Vec<String> = self
2157 .engine
2158 .open_positions()
2159 .iter()
2160 .map(|p| p.data.id.clone())
2161 .collect();
2162
2163 for id in open_ids {
2164 let symbol = match self.engine.get_position(&id) {
2165 Some(pos) => pos.data.symbol.clone(),
2166 None => continue,
2167 };
2168 let quote = match self.engine.last_quote(&symbol) {
2169 Some(q) => q.clone(),
2170 None => continue,
2171 };
2172
2173 if let Ok(effects) = self.engine.apply_action(
2174 Action::ClosePosition {
2175 position_id: id.clone(),
2176 },
2177 quote.ts,
2178 ) {
2179 self.executor
2180 .process_effects(&effects, &self.engine, "e);
2181 }
2182 }
2183 }
2184
2185 pub fn run_raw_signals_future_with_config<F: DataFeed>(
2192 self,
2193 source_feed: &mut F,
2194 raw_signals: Vec<RawSignal>,
2195 profile: Option<&ManagementProfile>,
2196 future: FutureQuoteConfig,
2197 ) -> BacktestResult {
2198 self.run_raw_signals_future_controlled(
2199 source_feed,
2200 raw_signals,
2201 profile,
2202 future,
2203 &mut || false,
2204 &mut |_| {},
2205 )
2206 .expect("non-cancellable replay cannot be cancelled")
2207 }
2208
2209 fn run_raw_signals_future_controlled<F, C, P>(
2210 mut self,
2211 source_feed: &mut F,
2212 raw_signals: Vec<RawSignal>,
2213 profile: Option<&ManagementProfile>,
2214 future: FutureQuoteConfig,
2215 is_cancelled: &mut C,
2216 on_progress: &mut P,
2217 ) -> std::result::Result<BacktestResult, ReplayCancelled>
2218 where
2219 F: DataFeed,
2220 C: FnMut() -> bool,
2221 P: FnMut(ReplayProgress),
2222 {
2223 self.future_config = Some(future.clone());
2224 if let Err(error) = validate_replay_config(&self.config, Some(&future), &raw_signals) {
2225 return Ok(rejected_future_result(
2226 &self.config,
2227 &future,
2228 self.evaluation_options,
2229 error,
2230 ));
2231 }
2232 if let Some(profile) = profile
2233 && let Err(error) = profile.validate()
2234 {
2235 return Ok(rejected_future_result(
2236 &self.config,
2237 &future,
2238 self.evaluation_options,
2239 error.to_string(),
2240 ));
2241 }
2242 if let Err(error) = self.validate_entry_profile_routes(&raw_signals) {
2243 return Ok(rejected_future_result(
2244 &self.config,
2245 &future,
2246 self.evaluation_options,
2247 error,
2248 ));
2249 }
2250
2251 let total_events = source_feed.total_events();
2252 let mut processed_events = 0;
2253 let mut ordered_events = Vec::<FeedEvent>::new();
2254 let mut source_last_ts = BTreeMap::<String, NaiveDateTime>::new();
2255 let mut invalid_quotes = 0_u64;
2256 while let Some(batch) = source_feed.next_batch() {
2257 if is_cancelled() {
2258 return Err(ReplayCancelled);
2259 }
2260 for feed_event in batch.events {
2261 if is_cancelled() {
2262 return Err(ReplayCancelled);
2263 }
2264 let symbol = feed_event.event.symbol().to_owned();
2265 let event_ts = feed_event.event.ts();
2266 if source_last_ts
2267 .get(&symbol)
2268 .is_some_and(|last| *last > event_ts)
2269 {
2270 invalid_quotes += 1;
2271 processed_events += 1;
2272 continue;
2273 }
2274 source_last_ts.insert(symbol, event_ts);
2275 ordered_events.push(feed_event);
2276 }
2277 }
2278 if is_cancelled() {
2279 return Err(ReplayCancelled);
2280 }
2281 ordered_events.sort_by_key(FeedEvent::ordering_key);
2282 if is_cancelled() {
2283 return Err(ReplayCancelled);
2284 }
2285 let primary_eod = ordered_events
2286 .iter()
2287 .filter(|event| event.metadata.roles.primary)
2288 .filter_map(|event| event.event.to_valid_quote())
2289 .map(|quote| quote.ts)
2290 .max();
2291 let mut ordered_feed = crate::data_feed::VecFeed::from_feed_events(ordered_events);
2292 let mut feed = DataFeedBatchAdapter {
2293 feed: &mut ordered_feed,
2294 };
2295 let mut hook = StaticReplayHook;
2296 match self.run_raw_signals_future_batches(
2297 &mut feed,
2298 primary_eod,
2299 raw_signals,
2300 profile,
2301 future,
2302 total_events,
2303 processed_events,
2304 invalid_quotes,
2305 is_cancelled,
2306 on_progress,
2307 &mut hook,
2308 ) {
2309 Ok(result) => Ok(result),
2310 Err(FutureBatchReplayError::Cancelled) => Err(ReplayCancelled),
2311 Err(FutureBatchReplayError::Feed(error)) => match error {},
2312 Err(FutureBatchReplayError::Dynamic) => {
2313 unreachable!("static replay has no dynamic hook")
2314 }
2315 }
2316 }
2317
2318 #[allow(clippy::too_many_arguments)]
2319 fn run_raw_signals_future_batches<F, C, P, H>(
2320 mut self,
2321 feed: &mut F,
2322 primary_eod: Option<NaiveDateTime>,
2323 raw_signals: Vec<RawSignal>,
2324 profile: Option<&ManagementProfile>,
2325 future: FutureQuoteConfig,
2326 known_total_events: Option<usize>,
2327 mut processed_events: usize,
2328 mut invalid_quotes: u64,
2329 is_cancelled: &mut C,
2330 on_progress: &mut P,
2331 hook: &mut H,
2332 ) -> std::result::Result<BacktestResult, FutureBatchReplayError<F::Error>>
2333 where
2334 F: FallibleBatchFeed,
2335 C: FnMut() -> bool,
2336 P: FnMut(ReplayProgress),
2337 H: FutureReplayHook,
2338 {
2339 let total_events = known_total_events.unwrap_or(0);
2340 let total_signals = raw_signals.len();
2341 let mut processed_signals = 0;
2342 on_progress(ReplayProgress {
2343 processed_events,
2344 total_events,
2345 processed_signals,
2346 total_signals,
2347 });
2348
2349 let execution_model = ExecutionModel::new(
2350 qs_core::types::ExecutionConvention::FutureQuoteV1,
2351 self.config.fill_model,
2352 if future.slippage_pips == 0.0 {
2353 SlippageModel::None
2354 } else {
2355 SlippageModel::FixedPips {
2356 pips: future.slippage_pips,
2357 }
2358 },
2359 );
2360 let pricer = ExecutionPricer::new(execution_model);
2361 let mut scheduled: Vec<ScheduledSignal> = raw_signals
2362 .into_iter()
2363 .enumerate()
2364 .map(|(sequence, signal)| {
2365 let signal_ts = signal.ts();
2366 ScheduledSignal::new(
2367 sequence as u64,
2368 signal_ts,
2369 signal_ts
2370 .checked_add_signed(Duration::milliseconds(future.signal_latency_ms))
2371 .expect("signal latency overflow was validated before scheduling"),
2372 signal,
2373 false,
2374 )
2375 })
2376 .collect();
2377 scheduled.sort_by_key(|signal| (signal.effective_ts, signal.sequence));
2378 let mut scheduled = VecDeque::from(scheduled);
2379 let mut queued = VecDeque::<QueuedAction>::new();
2380 let mut lifecycle = LifecycleLedger::new();
2381 let contract_sizes = effective_contract_sizes(&self.config);
2382 let mut future_executor = FutureExecutor::new(
2383 self.config.initial_balance,
2384 contract_sizes.clone(),
2385 future.pnl_epsilon,
2386 )
2387 .with_currency_plan(future.currency_plan.clone())
2388 .with_costs(
2389 self.config.costs.clone(),
2390 effective_point_sizes(&self.config),
2391 );
2392 let mut portfolio =
2393 PortfolioRecorder::new(self.config.initial_balance, contract_sizes.clone())
2394 .with_fill_model(self.config.fill_model)
2395 .with_stale_quote_after_millis(future.stale_quote_after_ms)
2396 .with_currency_plan(future.currency_plan.clone());
2397 let mut mtm_curve = MtmCurveCollector::new(future.mtm_output)
2398 .expect("MTM output policy was validated before replay");
2399 let mut last_mtm_candidate = None;
2400 let mut conversion_quotes =
2401 ConversionQuoteBook::new(Duration::milliseconds(future.conversion_stale_after_ms))
2402 .expect("conversion quote staleness was validated before replay");
2403 if let Some(plan) = future.currency_plan.as_ref() {
2404 for quote in plan.strict_before_warmup_quotes() {
2405 conversion_quotes
2406 .record_canonical_tick(quote.clone())
2407 .expect("currency plan warmups were validated during construction");
2408 }
2409 }
2410 let mut last_quote_ts = BTreeMap::<String, NaiveDateTime>::new();
2411 let mut unconverted_cost_events = 0u64;
2412 let mut zero_spread_bar_quotes = 0u64;
2413 let mut last_processed_primary_ts = None;
2414 let mut effective_terminal_ts = primary_eod;
2415 let mut terminated_quiescently = false;
2416 let bar_execution_timeframes = hook.bar_execution_timeframes();
2417
2418 while let Some(mut batch) =
2419 FallibleBatchFeed::next_batch(feed).map_err(FutureBatchReplayError::Feed)?
2420 {
2421 let batch_ts = batch.ts;
2422 if is_cancelled() {
2423 return Err(FutureBatchReplayError::Cancelled);
2424 }
2425 batch
2426 .events
2427 .sort_by_key(|event| (event.metadata.series_rank, event.metadata.row_sequence));
2428 let boundary_excursions = if hook.observes_position_economics() {
2429 portfolio.open_campaign_excursions()
2430 } else {
2431 BTreeMap::new()
2432 };
2433 let mut accepted = Vec::new();
2434 for feed_event in batch.events {
2435 if is_cancelled() {
2436 return Err(FutureBatchReplayError::Cancelled);
2437 }
2438 let fallback = self
2439 .config
2440 .bar_spread_fallback
2441 .get(feed_event.event.symbol())
2442 .copied();
2443 let available_at = feed_event.available_at();
2445 let delayed_bar = match &feed_event.event {
2446 MarketEvent::Bar {
2447 ts,
2448 timeframe_seconds: Some(seconds),
2449 ..
2450 } => ts
2451 .checked_add_signed(chrono::Duration::seconds(*seconds as i64))
2452 .is_some_and(|nominal_close| available_at > nominal_close),
2453 _ => false,
2454 };
2455 let prices = match feed_event.event.bar_execution_prices(fallback) {
2456 None => None,
2457 Some(prices) => match prices.executable() {
2458 Some(mut prices) => {
2459 prices.ts = available_at;
2460 Some(prices)
2461 }
2462 None => {
2463 invalid_quotes += 1;
2464 processed_events += 1;
2465 if should_report_progress(processed_events, total_events) {
2466 on_progress(ReplayProgress {
2467 processed_events,
2468 total_events,
2469 processed_signals,
2470 total_signals,
2471 });
2472 }
2473 continue;
2474 }
2475 },
2476 };
2477 let bar = prices.filter(|bar| {
2478 !delayed_bar
2479 && match (
2480 bar_execution_timeframes.get(&bar.symbol),
2481 bar.timeframe_seconds,
2482 ) {
2483 (Some(execution), Some(seconds)) => *execution == seconds,
2484 _ => true,
2485 }
2486 });
2487 let series_only =
2488 bar.is_none() && matches!(feed_event.event, MarketEvent::Bar { .. });
2489 let mut quote = match &bar {
2490 Some(bar) => bar.open_quote(),
2491 None => feed_event.event.to_quote_with_spread_fallback(fallback),
2492 };
2493 quote.ts = available_at;
2494 if bar.is_some() && quote.bid == quote.ask {
2495 zero_spread_bar_quotes += 1;
2496 }
2497 if ExecutionPricer::validate_quote("e).is_err()
2498 || last_quote_ts
2499 .get("e.symbol)
2500 .is_some_and(|last| *last > quote.ts)
2501 {
2502 invalid_quotes += 1;
2503 processed_events += 1;
2504 if should_report_progress(processed_events, total_events) {
2505 on_progress(ReplayProgress {
2506 processed_events,
2507 total_events,
2508 processed_signals,
2509 total_signals,
2510 });
2511 }
2512 continue;
2513 }
2514 last_quote_ts.insert(quote.symbol.clone(), quote.ts);
2515 accepted.push((feed_event, quote, bar, series_only));
2516 }
2517 let accepted_primary_events = accepted
2518 .iter()
2519 .filter(|(event, ..)| event.metadata.roles.primary)
2520 .map(|(event, ..)| event.clone())
2521 .collect::<Vec<_>>();
2522 if !hook.preflight_primary_events(&accepted_primary_events) {
2523 return Err(FutureBatchReplayError::Dynamic);
2524 }
2525
2526 for (feed_event, quote, ..) in &accepted {
2527 if is_cancelled() {
2528 return Err(FutureBatchReplayError::Cancelled);
2529 }
2530 if feed_event.metadata.roles.conversion
2531 && matches!(feed_event.event, MarketEvent::Tick { .. })
2532 && conversion_quotes
2533 .record_canonical_tick(quote.clone())
2534 .is_err()
2535 {
2536 invalid_quotes += 1;
2537 }
2538 }
2539 unconverted_cost_events += future_executor.charge_rollovers(
2541 batch_ts,
2542 &mut portfolio,
2543 Some(&conversion_quotes),
2544 );
2545 if let Some(state) = hook.portfolio_state() {
2546 state.begin_batch(batch_ts, future_executor.balance());
2547 }
2548
2549 let valuation_only = accepted
2550 .iter()
2551 .any(|event| event.0.metadata.roles.conversion)
2552 && !accepted.iter().any(|event| event.0.metadata.roles.primary)
2553 && primary_eod.is_some_and(|eod| batch_ts <= eod);
2554
2555 let mut primary_quotes = Vec::new();
2556 let mut primary_events = Vec::new();
2557 let mut batch_quotes = BTreeMap::new();
2558 let mut executed_bars = Vec::new();
2559 for (feed_event, quote, bar, series_only) in accepted {
2560 if is_cancelled() {
2561 return Err(FutureBatchReplayError::Cancelled);
2562 }
2563 if !feed_event.metadata.roles.primary || series_only {
2564 if feed_event.metadata.roles.primary {
2565 last_processed_primary_ts = Some(quote.ts);
2566 primary_events.push(feed_event);
2567 }
2568 processed_events += 1;
2569 if should_report_progress(processed_events, total_events) {
2570 on_progress(ReplayProgress {
2571 processed_events,
2572 total_events,
2573 processed_signals,
2574 total_signals,
2575 });
2576 }
2577 continue;
2578 }
2579
2580 last_processed_primary_ts = Some(quote.ts);
2581 primary_events.push(feed_event);
2582 if let Some(bar) = bar {
2583 executed_bars.push(bar);
2584 }
2585 portfolio.record_quote(quote.clone());
2586 if hook.output_ready() {
2587 observe_future_equity(
2588 &mut portfolio,
2589 &future_executor,
2590 quote.ts,
2591 &conversion_quotes,
2592 EquityObservationKind::PreSettlement,
2593 &mut mtm_curve,
2594 &mut last_mtm_candidate,
2595 false,
2596 );
2597 }
2598 batch_quotes.insert(quote.symbol.clone(), quote.clone());
2599 primary_quotes.push(quote);
2600 }
2601
2602 let bar_batch = hook.reads_completed_bars_only()
2604 && primary_events
2605 .iter()
2606 .any(|event| matches!(event.event, MarketEvent::Bar { .. }));
2607 let mut boundary_events = Some(primary_events);
2608 if bar_batch {
2609 let pre_events = if hook.retains_post_bar_boundary() {
2610 boundary_events
2611 .as_ref()
2612 .expect("boundary events are available")
2613 .clone()
2614 } else {
2615 boundary_events
2616 .take()
2617 .expect("boundary events are taken once")
2618 };
2619 self.run_future_boundary(
2620 hook,
2621 batch_ts,
2622 pre_events,
2623 &boundary_excursions,
2624 &primary_quotes,
2625 &batch_quotes,
2626 profile,
2627 &future,
2628 true,
2629 last_mtm_candidate
2630 .as_ref()
2631 .and_then(|point: &EquityPoint| point.drawdown_pct),
2632 &mut scheduled,
2633 &mut queued,
2634 &mut lifecycle,
2635 &mut future_executor,
2636 &mut portfolio,
2637 &pricer,
2638 &conversion_quotes,
2639 )?;
2640 }
2641
2642 if let Some(representative_quote) = primary_quotes.first() {
2643 let mut settled_quotes = vec![false; primary_quotes.len()];
2644
2645 self.execute_queued_future(
2647 &batch_quotes,
2648 false,
2649 &mut queued,
2650 &mut lifecycle,
2651 &mut future_executor,
2652 &mut portfolio,
2653 &pricer,
2654 &conversion_quotes,
2655 );
2656 let increasing_symbols =
2657 queued_exposure_symbols(&queued, &batch_quotes, representative_quote.ts);
2658 invalid_quotes += self.settle_future_batch_symbols(
2659 &primary_quotes,
2660 &mut settled_quotes,
2661 Some(&increasing_symbols),
2662 &mut lifecycle,
2663 &mut future_executor,
2664 &mut portfolio,
2665 &pricer,
2666 &conversion_quotes,
2667 );
2668 self.execute_queued_future(
2669 &batch_quotes,
2670 true,
2671 &mut queued,
2672 &mut lifecycle,
2673 &mut future_executor,
2674 &mut portfolio,
2675 &pricer,
2676 &conversion_quotes,
2677 );
2678
2679 while scheduled
2680 .front()
2681 .is_some_and(|signal| signal.effective_ts < representative_quote.ts)
2682 {
2683 if is_cancelled() {
2684 return Err(FutureBatchReplayError::Cancelled);
2685 }
2686 let signal = scheduled.pop_front().expect("front checked");
2687 self.schedule_future_signal(
2688 signal,
2689 profile,
2690 representative_quote,
2691 &batch_quotes,
2692 &mut queued,
2693 &mut lifecycle,
2694 &mut future_executor,
2695 &mut portfolio,
2696 &pricer,
2697 &conversion_quotes,
2698 );
2699 processed_signals += 1;
2700 if should_report_progress(processed_signals, total_signals) {
2701 on_progress(ReplayProgress {
2702 processed_events,
2703 total_events,
2704 processed_signals,
2705 total_signals,
2706 });
2707 }
2708
2709 self.execute_queued_future(
2710 &batch_quotes,
2711 false,
2712 &mut queued,
2713 &mut lifecycle,
2714 &mut future_executor,
2715 &mut portfolio,
2716 &pricer,
2717 &conversion_quotes,
2718 );
2719 let increasing_symbols =
2720 queued_exposure_symbols(&queued, &batch_quotes, representative_quote.ts);
2721 invalid_quotes += self.settle_future_batch_symbols(
2722 &primary_quotes,
2723 &mut settled_quotes,
2724 Some(&increasing_symbols),
2725 &mut lifecycle,
2726 &mut future_executor,
2727 &mut portfolio,
2728 &pricer,
2729 &conversion_quotes,
2730 );
2731 self.execute_queued_future(
2732 &batch_quotes,
2733 true,
2734 &mut queued,
2735 &mut lifecycle,
2736 &mut future_executor,
2737 &mut portfolio,
2738 &pricer,
2739 &conversion_quotes,
2740 );
2741 }
2742
2743 invalid_quotes += self.settle_future_batch_symbols(
2745 &primary_quotes,
2746 &mut settled_quotes,
2747 None,
2748 &mut lifecycle,
2749 &mut future_executor,
2750 &mut portfolio,
2751 &pricer,
2752 &conversion_quotes,
2753 );
2754
2755 while scheduled
2756 .front()
2757 .is_some_and(|signal| signal.effective_ts == representative_quote.ts)
2758 {
2759 if is_cancelled() {
2760 return Err(FutureBatchReplayError::Cancelled);
2761 }
2762 let signal = scheduled.pop_front().expect("front checked");
2763 self.schedule_future_signal(
2764 signal,
2765 profile,
2766 representative_quote,
2767 &batch_quotes,
2768 &mut queued,
2769 &mut lifecycle,
2770 &mut future_executor,
2771 &mut portfolio,
2772 &pricer,
2773 &conversion_quotes,
2774 );
2775 processed_signals += 1;
2776 if should_report_progress(processed_signals, total_signals) {
2777 on_progress(ReplayProgress {
2778 processed_events,
2779 total_events,
2780 processed_signals,
2781 total_signals,
2782 });
2783 }
2784 self.execute_queued_future(
2785 &batch_quotes,
2786 false,
2787 &mut queued,
2788 &mut lifecycle,
2789 &mut future_executor,
2790 &mut portfolio,
2791 &pricer,
2792 &conversion_quotes,
2793 );
2794 self.execute_queued_future(
2795 &batch_quotes,
2796 true,
2797 &mut queued,
2798 &mut lifecycle,
2799 &mut future_executor,
2800 &mut portfolio,
2801 &pricer,
2802 &conversion_quotes,
2803 );
2804 }
2805 }
2806
2807 for bar in &executed_bars {
2809 if is_cancelled() {
2810 return Err(FutureBatchReplayError::Cancelled);
2811 }
2812 invalid_quotes += self.settle_future_bar_range(
2813 bar,
2814 &mut lifecycle,
2815 &mut future_executor,
2816 &mut portfolio,
2817 &pricer,
2818 &conversion_quotes,
2819 );
2820 }
2821 for bar in &executed_bars {
2822 portfolio.record_quote(bar.close_quote());
2823 }
2824
2825 if let Some(events) = boundary_events.take() {
2826 self.run_future_boundary(
2827 hook,
2828 batch_ts,
2829 events,
2830 &boundary_excursions,
2831 &primary_quotes,
2832 &batch_quotes,
2833 profile,
2834 &future,
2835 false,
2836 last_mtm_candidate
2837 .as_ref()
2838 .and_then(|point: &EquityPoint| point.drawdown_pct),
2839 &mut scheduled,
2840 &mut queued,
2841 &mut lifecycle,
2842 &mut future_executor,
2843 &mut portfolio,
2844 &pricer,
2845 &conversion_quotes,
2846 )?;
2847 }
2848
2849 for quote in &primary_quotes {
2850 if !hook.output_ready() {
2851 processed_events += 1;
2852 continue;
2853 }
2854 observe_future_equity(
2855 &mut portfolio,
2856 &future_executor,
2857 quote.ts,
2858 &conversion_quotes,
2859 EquityObservationKind::PostOutput,
2860 &mut mtm_curve,
2861 &mut last_mtm_candidate,
2862 true,
2863 );
2864 processed_events += 1;
2865 if should_report_progress(processed_events, total_events) {
2866 on_progress(ReplayProgress {
2867 processed_events,
2868 total_events,
2869 processed_signals,
2870 total_signals,
2871 });
2872 }
2873 }
2874 if valuation_only && hook.output_ready() {
2875 observe_future_equity(
2876 &mut portfolio,
2877 &future_executor,
2878 batch_ts,
2879 &conversion_quotes,
2880 EquityObservationKind::ConversionRevaluation,
2881 &mut mtm_curve,
2882 &mut last_mtm_candidate,
2883 false,
2884 );
2885 }
2886 conversion_quotes.retain_replay_causal_predecessors(
2887 batch_ts,
2888 scheduled
2889 .iter()
2890 .map(|signal| signal.effective_ts)
2891 .chain(primary_eod),
2892 );
2893
2894 if !hook.is_active()
2895 && last_processed_primary_ts.is_some()
2896 && scheduled.is_empty()
2897 && queued.is_empty()
2898 && self.engine.open_positions().is_empty()
2899 && self.engine.pending_positions().is_empty()
2900 {
2901 effective_terminal_ts = last_processed_primary_ts;
2902 terminated_quiescently = true;
2903 break;
2904 }
2905 }
2906
2907 for action in queued {
2908 if is_cancelled() {
2909 return Err(FutureBatchReplayError::Cancelled);
2910 }
2911 if let Some(signal) = action.entry_signal.as_ref() {
2912 self.record_entry_resolution_rejection(
2913 action.action_id.clone(),
2914 signal,
2915 Some(&SelectedEntryProfile {
2916 profile: action.entry_profile.clone(),
2917 source: action
2918 .entry_profile_selection_source
2919 .unwrap_or(EntryProfileSelectionSource::Unprofiled),
2920 entry_class: match signal {
2921 RawSignal::Entry { entry_class, .. } => entry_class.clone(),
2922 _ => None,
2923 },
2924 profile_name: action.selected_profile_name.clone(),
2925 }),
2926 EntryResolutionStage::MarketExecution,
2927 None,
2928 None,
2929 "quote_eligibility",
2930 "no_eligible_quote".into(),
2931 );
2932 }
2933 let mut disposition =
2934 ActionDisposition::rejected(action.action_id, "no_eligible_quote");
2935 disposition.action_kind = Some(action.action_kind);
2936 disposition.signal_ts = Some(action.signal_ts);
2937 disposition.effective_ts = Some(action.effective_ts);
2938 self.record_disposition(&mut lifecycle, disposition);
2939 }
2940 for signal in scheduled {
2941 if is_cancelled() {
2942 return Err(FutureBatchReplayError::Cancelled);
2943 }
2944 let action_id = signal.resolved_action_id();
2945 if signal.signal.is_entry() {
2946 let selected = self
2947 .select_entry_profile(&signal.signal, profile, signal.instance)
2948 .ok();
2949 let stage = match &signal.signal {
2950 RawSignal::Entry {
2951 order_type: OrderType::Market,
2952 ..
2953 } => EntryResolutionStage::MarketExecution,
2954 _ => EntryResolutionStage::PendingPlacement,
2955 };
2956 self.record_entry_resolution_rejection(
2957 action_id.clone(),
2958 &signal.signal,
2959 selected.as_ref(),
2960 stage,
2961 None,
2962 None,
2963 "quote_eligibility",
2964 "no_eligible_quote".into(),
2965 );
2966 }
2967 let mut disposition = ActionDisposition::rejected(action_id, "no_eligible_quote");
2968 disposition.action_kind = Some(raw_signal_kind(&signal.signal).to_owned());
2969 disposition.signal_ts = Some(signal.signal_ts);
2970 disposition.effective_ts = Some(signal.effective_ts);
2971 self.record_disposition(&mut lifecycle, disposition);
2972 processed_signals += 1;
2973 if should_report_progress(processed_signals, total_signals) {
2974 on_progress(ReplayProgress {
2975 processed_events,
2976 total_events,
2977 processed_signals,
2978 total_signals,
2979 });
2980 }
2981 }
2982
2983 if is_cancelled() {
2984 return Err(FutureBatchReplayError::Cancelled);
2985 }
2986
2987 if self.config.close_on_finish {
2988 let execution_ts = effective_terminal_ts;
2989 let ids: Vec<String> = future_executor
2990 .open_snapshots()
2991 .into_iter()
2992 .map(|position| position.position_id)
2993 .collect();
2994 for (sequence, id) in ids.into_iter().enumerate() {
2995 if is_cancelled() {
2996 return Err(FutureBatchReplayError::Cancelled);
2997 }
2998 let Some(symbol) = self
2999 .engine
3000 .get_position(&id)
3001 .map(|position| position.data.symbol.clone())
3002 else {
3003 continue;
3004 };
3005 let Some(quote) = portfolio.quote(&symbol).cloned() else {
3006 continue;
3007 };
3008 let action_id = format!("end_of_data:{sequence:08}");
3009 let transaction = (|| -> Result<_, FutureTransactionError> {
3010 let side = self
3011 .engine
3012 .get_position(&id)
3013 .ok_or_else(|| {
3014 FutureApplyError::Core(qs_core::CoreError::PositionNotFound(id.clone()))
3015 })?
3016 .data
3017 .side;
3018 let execution = pricer
3019 .market_exit(side, "e, self.pip_size(&symbol))
3020 .map_err(FutureApplyError::from)?;
3021 let engine_transaction =
3022 self.engine.begin_close_position_with_reason_future_at(
3023 &id,
3024 CloseReason::EndOfData,
3025 "e,
3026 execution,
3027 execution_ts.unwrap_or(quote.ts),
3028 )?;
3029 let committed_effects = engine_transaction.effects().to_vec();
3030 let affected =
3031 if FutureExecutor::requires_processing(engine_transaction.effects()) {
3032 match future_executor.process_future_effects_with_currency(
3033 engine_transaction.effects(),
3034 &self.engine,
3035 "e,
3036 Some(&action_id),
3037 None,
3038 execution_ts.unwrap_or(quote.ts),
3039 &mut portfolio,
3040 Some(&conversion_quotes),
3041 ) {
3042 Ok(affected) => affected,
3043 Err(error) => {
3044 engine_transaction.rollback(&mut self.engine);
3045 return Err(error.into());
3046 }
3047 }
3048 } else {
3049 Vec::new()
3050 };
3051 let _ = engine_transaction.commit();
3052 self.record_committed_effects(committed_effects, Some(action_id.clone()));
3053 Ok(affected)
3054 })();
3055
3056 let mut disposition = match transaction {
3057 Ok(affected) => {
3058 let mut disposition = ActionDisposition::applied(action_id);
3059 disposition.position_ids = affected;
3060 disposition
3061 }
3062 Err(error) => ActionDisposition::failed(action_id, error.to_string()),
3063 };
3064 disposition.action_kind = Some("end_of_data".into());
3065 disposition.effective_ts = Some(execution_ts.unwrap_or(quote.ts));
3066 self.record_disposition(&mut lifecycle, disposition);
3067 }
3068 }
3069
3070 if is_cancelled() {
3071 return Err(FutureBatchReplayError::Cancelled);
3072 }
3073
3074 if let Some(ts) = effective_terminal_ts
3075 && hook.output_ready()
3076 {
3077 let observation_kind = if terminated_quiescently {
3078 EquityObservationKind::QuiescentTermination
3079 } else {
3080 EquityObservationKind::EndOfData
3081 };
3082 observe_future_equity(
3083 &mut portfolio,
3084 &future_executor,
3085 ts,
3086 &conversion_quotes,
3087 observation_kind,
3088 &mut mtm_curve,
3089 &mut last_mtm_candidate,
3090 false,
3091 );
3092 future_executor.finalize_pending_orders_at_end(ts);
3093 }
3094
3095 if !hook.on_final_committed(
3096 &mut self.committed_feedback,
3097 &mut self.committed_feedback_events,
3098 ) {
3099 return Err(FutureBatchReplayError::Dynamic);
3100 }
3101
3102 let pending_orders = self
3103 .engine
3104 .pending_positions()
3105 .into_iter()
3106 .map(|position| {
3107 let metadata = future_executor.pending_metadata(&position.data.id);
3108 PendingOrderSnapshot {
3109 position_id: position.data.id.clone(),
3110 action_id: metadata.as_ref().map(|value| value.0.clone()),
3111 signal_ts: metadata.as_ref().map(|value| value.1),
3112 effective_ts: metadata.as_ref().map(|value| value.2),
3113 symbol: position.data.symbol.clone(),
3114 side: position.data.side,
3115 order_type: position.data.order_type,
3116 requested_price: position.data.pending_price,
3117 size: position.data.size,
3118 initial_stop: position.current_stoploss(),
3119 group: position.data.group.clone(),
3120 trade_id: position.data.trade_id.clone(),
3121 }
3122 })
3123 .collect();
3124 let mut tags = BTreeMap::new();
3125 tags.insert("invalid_quote_count".into(), invalid_quotes.to_string());
3126 tags.insert(
3127 "termination_reason".into(),
3128 if terminated_quiescently {
3129 "quiescent"
3130 } else {
3131 "end_of_data"
3132 }
3133 .into(),
3134 );
3135 insert_economic_support_metadata(&mut tags, &self.config);
3136 let (equity_curve, mtm_output_summary) = mtm_curve.into_parts();
3137 let entry_profile_default = self
3138 .entry_profiles
3139 .as_ref()
3140 .and_then(|profiles| profiles.default_profile().cloned());
3141 let entry_profile_routes = self
3142 .entry_profiles
3143 .as_ref()
3144 .map(|profiles| profiles.routes().clone())
3145 .unwrap_or_default();
3146 let artifacts = FutureBacktestArtifacts {
3147 format_version: FUTURE_ARTIFACT_FORMAT_VERSION,
3148 execution: ExecutionMetadata {
3149 execution_model,
3150 initial_balance: self.config.initial_balance,
3151 account_currency: future
3152 .currency_plan
3153 .as_ref()
3154 .map(|plan| plan.account_currency().to_owned()),
3155 currency_plan: future.currency_plan.clone(),
3156 contract_sizes: contract_sizes.into_iter().collect(),
3157 instrument_manifest: self.config.instrument_manifest.clone(),
3158 instrument_sizing: std::mem::take(&mut self.instrument_sizing),
3159 market_entry_sizing_basis: future.market_entry_sizing_basis,
3160 market_entry_sizing: std::mem::take(&mut self.market_entry_sizing),
3161 entry_profile_default,
3162 entry_profile_routes,
3163 entry_profile_resolutions: std::mem::take(&mut self.entry_profile_resolutions),
3164 costs: self.config.costs.clone().into_iter().collect(),
3165 unconverted_cost_events,
3166 zero_spread_bar_quotes,
3167 stale_quote_after_millis: future.stale_quote_after_ms,
3168 pnl_epsilon: future.pnl_epsilon,
3169 tags,
3170 run_tags: self.config.run_tags.clone(),
3171 position_tags: hook
3172 .portfolio_state()
3173 .map(|state| state.position_tags(&future_executor.fills))
3174 .unwrap_or_default(),
3175 ..ExecutionMetadata::default()
3176 },
3177 fills: future_executor.fills.clone(),
3178 close_events: future_executor.close_events.clone(),
3179 cost_events: future_executor.cost_events.clone(),
3180 completed_positions: future_executor.completed_positions.clone(),
3181 open_positions: portfolio.latest_open_positions().to_vec(),
3182 pending_orders,
3183 pending_order_lifecycle: future_executor.pending_order_lifecycle,
3184 lifecycle,
3185 equity_curve,
3186 mtm_output_summary,
3187 max_drawdown: portfolio.max_drawdown(),
3188 max_drawdown_pct: portfolio.max_drawdown_pct(),
3189 };
3190 on_progress(ReplayProgress {
3191 processed_events,
3192 total_events: known_total_events.unwrap_or(processed_events),
3193 processed_signals,
3194 total_signals,
3195 });
3196 Ok(BacktestResult::from_future_artifacts_with_options(
3197 artifacts,
3198 self.evaluation_options,
3199 ))
3200 }
3201
3202 #[allow(clippy::too_many_arguments)]
3206 fn run_future_boundary<H, E>(
3207 &mut self,
3208 hook: &mut H,
3209 batch_ts: NaiveDateTime,
3210 primary_events: Vec<FeedEvent>,
3211 boundary_excursions: &BTreeMap<String, crate::portfolio::CampaignExcursion>,
3212 primary_quotes: &[PriceQuote],
3213 batch_quotes: &BTreeMap<String, PriceQuote>,
3214 profile: Option<&ManagementProfile>,
3215 future: &FutureQuoteConfig,
3216 decided_before_quotes: bool,
3217 drawdown_fraction: Option<f64>,
3218 scheduled: &mut VecDeque<ScheduledSignal>,
3219 queued: &mut VecDeque<QueuedAction>,
3220 lifecycle: &mut LifecycleLedger,
3221 future_executor: &mut FutureExecutor,
3222 portfolio: &mut PortfolioRecorder,
3223 pricer: &ExecutionPricer,
3224 conversion_quotes: &ConversionQuoteBook,
3225 ) -> std::result::Result<(), FutureBatchReplayError<E>>
3226 where
3227 H: FutureReplayHook,
3228 {
3229 let strategy_batch = TimestampBatch {
3230 ts: batch_ts,
3231 events: primary_events,
3232 };
3233 let generated = {
3234 let positions = BoundaryPositionFacts::new(boundary_excursions, future_executor);
3235 if decided_before_quotes {
3236 hook.on_pre_bar_boundary(
3237 &strategy_batch,
3238 &self.engine,
3239 lifecycle,
3240 &positions,
3241 &mut self.committed_feedback,
3242 &mut self.committed_feedback_events,
3243 )
3244 } else {
3245 hook.on_boundary(
3246 &strategy_batch,
3247 &self.engine,
3248 lifecycle,
3249 &positions,
3250 &mut self.committed_feedback,
3251 &mut self.committed_feedback_events,
3252 )
3253 }
3254 .ok_or(FutureBatchReplayError::Dynamic)?
3255 };
3256 if !generated.is_empty() {
3257 let instances = generated
3258 .iter()
3259 .map(|scheduled| scheduled.instance)
3260 .collect::<BTreeSet<_>>();
3261 for instance in instances {
3262 let generated_signals = generated
3263 .iter()
3264 .filter(|scheduled| scheduled.instance == instance)
3265 .map(|scheduled| scheduled.signal.clone())
3266 .collect::<Vec<_>>();
3267 if let Err(error) =
3268 validate_replay_config(&self.config, Some(future), &generated_signals)
3269 {
3270 hook.reject_generated_configuration(instance, error);
3271 return Err(FutureBatchReplayError::Dynamic);
3272 }
3273 }
3274 }
3275 let generated = self.supervise_generated(
3276 hook,
3277 batch_ts,
3278 generated,
3279 decided_before_quotes,
3280 drawdown_fraction,
3281 scheduled,
3282 queued,
3283 lifecycle,
3284 future_executor,
3285 );
3286 for mut generated_signal in generated {
3287 if decided_before_quotes {
3288 generated_signal.requires_later_quote = false;
3289 }
3290 if generated_signal.effective_ts <= batch_ts
3291 && let Some(representative_quote) = primary_quotes.first()
3292 {
3293 self.schedule_future_signal(
3294 generated_signal,
3295 profile,
3296 representative_quote,
3297 batch_quotes,
3298 queued,
3299 lifecycle,
3300 future_executor,
3301 portfolio,
3302 pricer,
3303 conversion_quotes,
3304 );
3305 } else {
3306 scheduled.push_back(generated_signal);
3307 }
3308 }
3309 Ok(())
3310 }
3311
3312 #[allow(clippy::too_many_arguments)]
3316 fn settle_future_bar_range(
3317 &mut self,
3318 bar: &BarExecutionPrices,
3319 lifecycle: &mut LifecycleLedger,
3320 future_executor: &mut FutureExecutor,
3321 portfolio: &mut PortfolioRecorder,
3322 pricer: &ExecutionPricer,
3323 conversion_quotes: &ConversionQuoteBook,
3324 ) -> u64 {
3325 let mut failures = 0;
3326 for side in [Side::Buy, Side::Sell] {
3327 let legs = match side {
3328 Side::Buy => [(bar.open, bar.low), (bar.low, bar.high)],
3329 Side::Sell => [(bar.open, bar.high), (bar.high, bar.low)],
3330 };
3331 for (from, to) in legs {
3332 failures += self.walk_future_bar_leg(
3333 bar,
3334 side,
3335 from,
3336 to,
3337 lifecycle,
3338 future_executor,
3339 portfolio,
3340 pricer,
3341 conversion_quotes,
3342 );
3343 }
3344 }
3345 failures
3346 }
3347
3348 #[allow(clippy::too_many_arguments)]
3349 fn walk_future_bar_leg(
3350 &mut self,
3351 bar: &BarExecutionPrices,
3352 side: Side,
3353 from: f64,
3354 to: f64,
3355 lifecycle: &mut LifecycleLedger,
3356 future_executor: &mut FutureExecutor,
3357 portfolio: &mut PortfolioRecorder,
3358 pricer: &ExecutionPricer,
3359 conversion_quotes: &ConversionQuoteBook,
3360 ) -> u64 {
3361 if !(from.is_finite() && to.is_finite()) || from == to {
3362 return 0;
3363 }
3364 let rising = to > from;
3365 let ahead = |mid: f64, cursor: f64| {
3366 if rising {
3367 mid > cursor && mid <= to
3368 } else {
3369 mid < cursor && mid >= to
3370 }
3371 };
3372 let mut failures = 0;
3373 let mut cursor = from;
3374 for _ in 0..MAX_BAR_LEG_STEPS {
3376 let next = self
3377 .bar_trigger_levels(&bar.symbol, side, bar.half_spread)
3378 .into_iter()
3379 .filter(|(mid, _)| ahead(*mid, cursor))
3380 .min_by(|left, right| {
3381 let order = left.0.total_cmp(&right.0);
3382 if rising { order } else { order.reverse() }
3383 });
3384 let Some((mid, quote_prices)) = next else {
3385 break;
3386 };
3387 cursor = mid;
3388 let quote = PriceQuote {
3389 symbol: bar.symbol.clone(),
3390 ts: bar.ts,
3391 bid: quote_prices.0,
3392 ask: quote_prices.1,
3393 };
3394 if ExecutionPricer::validate_quote("e).is_err() {
3395 continue;
3396 }
3397 if self
3398 .settle_future_quote(
3399 "e,
3400 Some(side),
3401 lifecycle,
3402 future_executor,
3403 portfolio,
3404 pricer,
3405 conversion_quotes,
3406 )
3407 .is_err()
3408 {
3409 failures += 1;
3410 }
3411 }
3412 let extreme = bar.quote_at_mid(to);
3413 if ExecutionPricer::validate_quote(&extreme).is_ok()
3414 && self
3415 .settle_future_quote(
3416 &extreme,
3417 Some(side),
3418 lifecycle,
3419 future_executor,
3420 portfolio,
3421 pricer,
3422 conversion_quotes,
3423 )
3424 .is_err()
3425 {
3426 failures += 1;
3427 }
3428 failures
3429 }
3430
3431 fn bar_trigger_levels(
3433 &self,
3434 symbol: &str,
3435 side: Side,
3436 half_spread: f64,
3437 ) -> Vec<(f64, (f64, f64))> {
3438 let model = self.config.fill_model;
3439 let spread = half_spread * 2.0;
3440 let mut levels = Vec::new();
3441 let ids = self
3442 .engine
3443 .manager
3444 .pending_ids_by_symbol_sorted(symbol)
3445 .into_iter()
3446 .chain(self.engine.manager.open_ids_by_symbol_sorted(symbol));
3447 for id in ids {
3448 let Some(position) = self.engine.get_position(&id) else {
3449 continue;
3450 };
3451 if position.data.side != side {
3452 continue;
3453 }
3454 let compared = match (position.data.status, model, side) {
3456 (_, FillModel::MidPrice, _) => QuoteSide::Mid,
3457 (_, FillModel::AskOnly, _) => QuoteSide::Ask,
3458 (PositionStatus::Pending, FillModel::BidAsk, Side::Buy)
3459 | (PositionStatus::Open, FillModel::BidAsk, Side::Sell) => QuoteSide::Ask,
3460 _ => QuoteSide::Bid,
3461 };
3462 for level in position.future_trigger_levels() {
3463 levels.push(match compared {
3464 QuoteSide::Bid => (level + half_spread, (level, level + spread)),
3465 QuoteSide::Ask => (level - half_spread, (level - spread, level)),
3466 QuoteSide::Mid => {
3468 let (bid, ask) = (level - half_spread, level + half_spread);
3469 if (bid + ask) / 2.0 == level {
3470 (level, (bid, ask))
3471 } else {
3472 (level, (level, level))
3473 }
3474 }
3475 });
3476 }
3477 }
3478 levels
3479 }
3480
3481 #[allow(clippy::too_many_arguments)]
3482 fn settle_future_batch_symbols(
3483 &mut self,
3484 primary_quotes: &[PriceQuote],
3485 settled_quotes: &mut [bool],
3486 symbols: Option<&BTreeSet<String>>,
3487 lifecycle: &mut LifecycleLedger,
3488 future_executor: &mut FutureExecutor,
3489 portfolio: &mut PortfolioRecorder,
3490 pricer: &ExecutionPricer,
3491 conversion_quotes: &ConversionQuoteBook,
3492 ) -> u64 {
3493 let mut failures = 0;
3494 for (index, quote) in primary_quotes.iter().enumerate() {
3495 if settled_quotes[index]
3496 || symbols.is_some_and(|symbols| !symbols.contains("e.symbol))
3497 {
3498 continue;
3499 }
3500 if self
3501 .settle_future_quote(
3502 quote,
3503 None,
3504 lifecycle,
3505 future_executor,
3506 portfolio,
3507 pricer,
3508 conversion_quotes,
3509 )
3510 .is_err()
3511 {
3512 failures += 1;
3513 }
3514 settled_quotes[index] = true;
3515 }
3516 failures
3517 }
3518
3519 #[allow(clippy::too_many_arguments)]
3520 fn settle_future_quote(
3521 &mut self,
3522 quote: &PriceQuote,
3523 side: Option<Side>,
3524 lifecycle: &mut LifecycleLedger,
3525 future_executor: &mut FutureExecutor,
3526 portfolio: &mut PortfolioRecorder,
3527 pricer: &ExecutionPricer,
3528 conversion_quotes: &ConversionQuoteBook,
3529 ) -> Result<(), FutureTransactionError> {
3530 let (prepared, failures) = self.prepare_triggering_pending(quote, side, pricer);
3531 for (position_id, error) in failures {
3532 let action_id = format!("pending_execution:{position_id}");
3533 let mut disposition = ActionDisposition::rejected(action_id.clone(), error);
3534 disposition.action_kind = Some("pending_execution".into());
3535 disposition.effective_ts = Some(quote.ts);
3536 disposition.position_ids.push(position_id.clone());
3537 self.record_disposition(lifecycle, disposition);
3538 if let Ok(engine_transaction) = self.engine.begin_future_action(
3539 Action::CancelPending {
3540 position_id: position_id.clone(),
3541 },
3542 quote.ts,
3543 ) {
3544 let committed_effects = engine_transaction.effects().to_vec();
3545 if FutureExecutor::requires_processing(engine_transaction.effects())
3546 && let Err(error) = future_executor.process_future_effects_with_currency(
3547 engine_transaction.effects(),
3548 &self.engine,
3549 quote,
3550 Some(&action_id),
3551 None,
3552 quote.ts,
3553 portfolio,
3554 Some(conversion_quotes),
3555 )
3556 {
3557 engine_transaction.rollback(&mut self.engine);
3558 return Err(error.into());
3559 }
3560 let _ = engine_transaction.commit();
3561 self.record_committed_effects(committed_effects, Some(action_id.clone()));
3562 }
3563 }
3564
3565 let pending_action_ids = prepared
3566 .iter()
3567 .filter_map(|pending| {
3568 future_executor
3569 .pending_metadata(&pending.position_id)
3570 .map(|metadata| (pending.position_id.clone(), metadata.0))
3571 })
3572 .collect::<BTreeMap<_, _>>();
3573 let pip_size = self.pip_size("e.symbol);
3574 let engine_transaction = match side {
3575 Some(side) => self.engine.begin_on_price_future_effects_priced_for_side(
3576 quote, &prepared, pricer, pip_size, side,
3577 )?,
3578 None => self
3579 .engine
3580 .begin_on_price_future_effects_priced(quote, &prepared, pricer, pip_size)?,
3581 };
3582 let committed_effects = engine_transaction.effects().to_vec();
3583 if FutureExecutor::requires_processing(engine_transaction.effects())
3584 && let Err(error) = future_executor.process_future_effects_with_currency(
3585 engine_transaction.effects(),
3586 &self.engine,
3587 quote,
3588 None,
3589 None,
3590 quote.ts,
3591 portfolio,
3592 Some(conversion_quotes),
3593 )
3594 {
3595 engine_transaction.rollback(&mut self.engine);
3596 return Err(error.into());
3597 }
3598 let _ = engine_transaction.commit();
3599 for effect in committed_effects {
3600 let action_id = pending_fill_position_id(&effect)
3601 .and_then(|position_id| pending_action_ids.get(position_id))
3602 .cloned();
3603 self.record_committed_effects(vec![effect], action_id);
3604 }
3605 Ok(())
3606 }
3607
3608 #[allow(clippy::too_many_arguments)]
3609 fn schedule_future_signal(
3610 &mut self,
3611 scheduled: ScheduledSignal,
3612 profile: Option<&ManagementProfile>,
3613 quote: &PriceQuote,
3614 batch_quotes: &BTreeMap<String, PriceQuote>,
3615 queued: &mut VecDeque<QueuedAction>,
3616 lifecycle: &mut LifecycleLedger,
3617 future_executor: &mut FutureExecutor,
3618 portfolio: &mut PortfolioRecorder,
3619 pricer: &ExecutionPricer,
3620 conversion_quotes: &ConversionQuoteBook,
3621 ) {
3622 let explicit_action_id = scheduled.explicit_action_id;
3623 let base_id = scheduled.resolved_action_id();
3624 let selected_profile = if scheduled.signal.is_entry() {
3625 match self.select_entry_profile(&scheduled.signal, profile, scheduled.instance) {
3626 Ok(profile) => profile,
3627 Err(error) => {
3628 let reason = error.to_string();
3629 let stage = match scheduled.signal {
3630 RawSignal::Entry {
3631 order_type: OrderType::Market,
3632 ..
3633 } => EntryResolutionStage::MarketExecution,
3634 _ => EntryResolutionStage::PendingPlacement,
3635 };
3636 self.record_entry_resolution_rejection(
3637 base_id.clone(),
3638 &scheduled.signal,
3639 None,
3640 stage,
3641 None,
3642 None,
3643 "profile_selection",
3644 reason.clone(),
3645 );
3646 let mut disposition = ActionDisposition::rejected(base_id, reason);
3647 disposition.action_kind = Some("entry".into());
3648 disposition.signal_ts = Some(scheduled.signal_ts);
3649 disposition.effective_ts = Some(scheduled.effective_ts);
3650 self.record_disposition(lifecycle, disposition);
3651 return;
3652 }
3653 }
3654 } else {
3655 SelectedEntryProfile {
3656 profile: None,
3657 source: EntryProfileSelectionSource::Unprofiled,
3658 entry_class: None,
3659 profile_name: None,
3660 }
3661 };
3662 if let RawSignal::Entry {
3663 symbol,
3664 side,
3665 order_type: OrderType::Market,
3666 ..
3667 } = &scheduled.signal
3668 {
3669 queued.push_back(QueuedAction {
3670 action_id: base_id,
3671 action_kind: "entry".into(),
3672 action: Action::Open {
3673 symbol: symbol.clone(),
3674 side: *side,
3675 order_type: OrderType::Market,
3676 price: None,
3677 size: 1.0,
3678 stoploss: None,
3679 targets: Vec::new(),
3680 rules: Vec::new(),
3681 group: None,
3682 trade_id: None,
3683 },
3684 execution: None,
3685 symbol: symbol.clone(),
3686 signal_ts: scheduled.signal_ts,
3687 effective_ts: scheduled.effective_ts,
3688 entry_signal: Some(scheduled.signal),
3689 entry_profile: selected_profile.profile,
3690 entry_profile_selection_source: Some(selected_profile.source),
3691 selected_profile_name: selected_profile.profile_name,
3692 market_entry_sizing_audit: None,
3693 entry_profile_resolution_audit: None,
3694 requires_later_quote: scheduled.requires_later_quote,
3695 });
3696 return;
3697 }
3698 if scheduled.signal.is_entry() {
3699 let resolved = match selected_profile.profile.as_ref() {
3700 Some(profile) => self.resolve_profiled_entry(profile, &scheduled.signal),
3701 None => resolve_unprofiled_entry(&scheduled.signal),
3702 };
3703 match resolved {
3704 Ok(Some(resolved)) => {
3705 let entry_quote = batch_quotes.get(&resolved.symbol).unwrap_or(quote);
3706 self.enqueue_resolved_entry(
3707 base_id,
3708 scheduled,
3709 resolved,
3710 entry_quote,
3711 lifecycle,
3712 future_executor,
3713 portfolio,
3714 pricer,
3715 conversion_quotes,
3716 &selected_profile,
3717 )
3718 }
3719 Ok(None) => {
3720 let mut disposition = ActionDisposition::skipped(base_id, "not_an_entry");
3721 disposition.action_kind = Some("entry".into());
3722 disposition.signal_ts = Some(scheduled.signal_ts);
3723 disposition.effective_ts = Some(scheduled.effective_ts);
3724 self.record_disposition(lifecycle, disposition);
3725 }
3726 Err(error) => {
3727 let reason = error.to_string();
3728 self.record_entry_resolution_rejection(
3729 base_id.clone(),
3730 &scheduled.signal,
3731 Some(&selected_profile),
3732 EntryResolutionStage::PendingPlacement,
3733 match &scheduled.signal {
3734 RawSignal::Entry { price, .. } => *price,
3735 _ => None,
3736 },
3737 None,
3738 "profile_resolution",
3739 reason.clone(),
3740 );
3741 let mut disposition = ActionDisposition::rejected(base_id, reason);
3742 disposition.action_kind = Some("entry".into());
3743 disposition.signal_ts = Some(scheduled.signal_ts);
3744 disposition.effective_ts = Some(scheduled.effective_ts);
3745 self.record_disposition(lifecycle, disposition);
3746 }
3747 }
3748 return;
3749 }
3750
3751 let actions = self.resolve_future_actions(&scheduled.signal);
3752 if actions.is_empty() {
3753 let mut disposition = ActionDisposition::skipped(base_id, "position_not_found");
3754 disposition.action_kind = Some(raw_signal_kind(&scheduled.signal).to_owned());
3755 disposition.signal_ts = Some(scheduled.signal_ts);
3756 disposition.effective_ts = Some(scheduled.effective_ts);
3757 self.record_disposition(lifecycle, disposition);
3758 return;
3759 }
3760 let action_count = actions.len();
3761 if explicit_action_id && action_count != 1 {
3762 let mut disposition = ActionDisposition::rejected(
3763 base_id,
3764 "configured_command_resolved_multiple_actions",
3765 );
3766 disposition.action_kind = Some(raw_signal_kind(&scheduled.signal).to_owned());
3767 disposition.signal_ts = Some(scheduled.signal_ts);
3768 disposition.effective_ts = Some(scheduled.effective_ts);
3769 self.record_disposition(lifecycle, disposition);
3770 return;
3771 }
3772 for (index, action) in actions.into_iter().enumerate() {
3773 let action_id = if explicit_action_id && action_count == 1 {
3774 base_id.clone()
3775 } else {
3776 format!("{base_id}:action:{index:03}")
3777 };
3778 let Some(symbol) = self.action_symbol(&action) else {
3779 self.apply_future_action(
3780 action_id,
3781 raw_signal_kind(&scheduled.signal).to_owned(),
3782 action,
3783 None,
3784 scheduled.signal_ts,
3785 scheduled.effective_ts,
3786 quote,
3787 lifecycle,
3788 future_executor,
3789 portfolio,
3790 pricer,
3791 conversion_quotes,
3792 );
3793 continue;
3794 };
3795 if is_fill_bearing(&action) {
3796 queued.push_back(QueuedAction {
3797 action_id,
3798 action_kind: raw_signal_kind(&scheduled.signal).to_owned(),
3799 action,
3800 execution: None,
3801 symbol,
3802 signal_ts: scheduled.signal_ts,
3803 effective_ts: scheduled.effective_ts,
3804 entry_signal: None,
3805 entry_profile: None,
3806 entry_profile_selection_source: None,
3807 selected_profile_name: None,
3808 market_entry_sizing_audit: None,
3809 entry_profile_resolution_audit: None,
3810 requires_later_quote: scheduled.requires_later_quote,
3811 });
3812 } else {
3813 let action_quote = batch_quotes.get(&symbol).unwrap_or(quote);
3814 self.apply_future_action(
3815 action_id,
3816 raw_signal_kind(&scheduled.signal).to_owned(),
3817 action,
3818 None,
3819 scheduled.signal_ts,
3820 scheduled.effective_ts,
3821 action_quote,
3822 lifecycle,
3823 future_executor,
3824 portfolio,
3825 pricer,
3826 conversion_quotes,
3827 );
3828 }
3829 }
3830 }
3831
3832 #[allow(clippy::too_many_arguments)]
3833 fn enqueue_resolved_entry(
3834 &mut self,
3835 action_id: String,
3836 scheduled: ScheduledSignal,
3837 resolved: ResolvedEntry,
3838 quote: &PriceQuote,
3839 lifecycle: &mut LifecycleLedger,
3840 future_executor: &mut FutureExecutor,
3841 portfolio: &mut PortfolioRecorder,
3842 pricer: &ExecutionPricer,
3843 conversion_quotes: &ConversionQuoteBook,
3844 selected_profile: &SelectedEntryProfile,
3845 ) {
3846 let level_reference_price = resolved.price;
3847 let level_resolution = resolved.level_resolution.clone();
3848 match self.finalize_resolved_entry(
3849 resolved,
3850 future_executor.balance(),
3851 scheduled.effective_ts,
3852 Some(conversion_quotes),
3853 None,
3854 ) {
3855 Ok(finalized) => {
3856 let original_signal_price = match &scheduled.signal {
3857 RawSignal::Entry { price, .. } => *price,
3858 _ => None,
3859 };
3860 let (trade_id, level_reference_price) = match &finalized.action {
3861 Action::Open {
3862 trade_id, price, ..
3863 } => (
3864 trade_id.clone(),
3865 price.unwrap_or(quote.open_price(match &finalized.action {
3866 Action::Open { side, .. } => *side,
3867 _ => unreachable!(),
3868 })),
3869 ),
3870 _ => unreachable!("finalized entry must be an open action"),
3871 };
3872 let mut audit = EntryProfileResolutionAudit {
3873 action_id: action_id.clone(),
3874 trade_id,
3875 entry_class: selected_profile.entry_class.clone(),
3876 selection_source: selected_profile.source,
3877 selected_profile_name: selected_profile.profile_name.clone(),
3878 resolution_stage: EntryResolutionStage::PendingPlacement,
3879 original_signal_price,
3880 level_reference_price: Some(level_reference_price),
3881 level_resolution: Some(finalized.level_resolution.clone()),
3882 target_resolution: Some(finalized.target_resolution.clone()),
3883 configured_weights: finalized.configured_weights.clone(),
3884 allocated_target_steps: finalized.allocated_target_steps.clone(),
3885 remainder_steps: finalized.remainder_steps,
3886 outcome: crate::ledger::ActionDispositionStatus::Applied,
3887 rejection_stage: None,
3888 reason: None,
3889 };
3890 let committed = self.apply_future_action(
3891 action_id,
3892 "entry".into(),
3893 finalized.action,
3894 None,
3895 scheduled.signal_ts,
3896 scheduled.effective_ts,
3897 quote,
3898 lifecycle,
3899 future_executor,
3900 portfolio,
3901 pricer,
3902 conversion_quotes,
3903 );
3904 if committed {
3905 self.entry_profile_resolutions.push(audit);
3906 } else {
3907 if let Some(disposition) = lifecycle
3908 .as_slice()
3909 .iter()
3910 .rev()
3911 .find(|disposition| disposition.action_id == audit.action_id)
3912 {
3913 audit.outcome = disposition.status;
3914 audit.rejection_stage = Some("engine_or_accounting".into());
3915 audit.reason = disposition.reason.clone();
3916 }
3917 self.entry_profile_resolutions.push(audit);
3918 }
3919 }
3920 Err(error) => {
3921 self.record_entry_resolution_rejection(
3922 action_id.clone(),
3923 &scheduled.signal,
3924 Some(selected_profile),
3925 EntryResolutionStage::PendingPlacement,
3926 level_reference_price,
3927 Some(level_resolution),
3928 "sizing",
3929 error.clone(),
3930 );
3931 let mut disposition = ActionDisposition::rejected(action_id, error);
3932 disposition.action_kind = Some("entry".into());
3933 disposition.signal_ts = Some(scheduled.signal_ts);
3934 disposition.effective_ts = Some(scheduled.effective_ts);
3935 self.record_disposition(lifecycle, disposition);
3936 }
3937 }
3938 }
3939
3940 #[allow(clippy::too_many_arguments)]
3941 fn execute_queued_future(
3942 &mut self,
3943 quotes: &BTreeMap<String, PriceQuote>,
3944 exposure_increasing: bool,
3945 queued: &mut VecDeque<QueuedAction>,
3946 lifecycle: &mut LifecycleLedger,
3947 future_executor: &mut FutureExecutor,
3948 portfolio: &mut PortfolioRecorder,
3949 pricer: &ExecutionPricer,
3950 conversion_quotes: &ConversionQuoteBook,
3951 ) {
3952 let mut remaining = VecDeque::new();
3953 while let Some(mut action) = queued.pop_front() {
3954 let increases = is_exposure_increasing(&action.action);
3955 let Some(quote) = quotes.get(&action.symbol) else {
3956 remaining.push_back(action);
3957 continue;
3958 };
3959 if action.effective_ts > quote.ts
3960 || (action.requires_later_quote && quote.ts <= action.signal_ts)
3961 || increases != exposure_increasing
3962 {
3963 remaining.push_back(action);
3964 continue;
3965 }
3966
3967 if let Some(mut signal) = action.entry_signal.take() {
3968 let (side, symbol, original_signal_price, entry_class) = match &signal {
3969 RawSignal::Entry {
3970 side,
3971 symbol,
3972 price,
3973 entry_class,
3974 ..
3975 } => (*side, symbol.clone(), *price, entry_class.clone()),
3976 _ => unreachable!("queued entry metadata must contain an entry signal"),
3977 };
3978 let execution = match pricer.market_entry(side, quote, self.pip_size(&symbol)) {
3979 Ok(fill) => fill,
3980 Err(error) => {
3981 let mut disposition =
3982 ActionDisposition::rejected(action.action_id, error.to_string());
3983 disposition.action_kind = Some(action.action_kind);
3984 disposition.signal_ts = Some(action.signal_ts);
3985 disposition.effective_ts = Some(action.effective_ts);
3986 self.record_disposition(lifecycle, disposition);
3987 continue;
3988 }
3989 };
3990 if let RawSignal::Entry { price, .. } = &mut signal {
3991 *price = Some(execution.price);
3992 }
3993 action.execution = Some(execution);
3994 let configured_basis = self
3995 .future_config
3996 .as_ref()
3997 .map(|config| config.market_entry_sizing_basis)
3998 .unwrap_or_default();
3999 let (applied_basis, fallback_to_fill, sizing_reference_price) =
4000 match (configured_basis, original_signal_price) {
4001 (MarketEntrySizingBasis::SignalEntryPrice, Some(price)) => {
4002 (MarketEntrySizingBasis::SignalEntryPrice, false, price)
4003 }
4004 (MarketEntrySizingBasis::SignalEntryPrice, None) => {
4005 (MarketEntrySizingBasis::FillPrice, true, execution.price)
4006 }
4007 (MarketEntrySizingBasis::FillPrice, _) => {
4008 (MarketEntrySizingBasis::FillPrice, false, execution.price)
4009 }
4010 };
4011 let resolved = match action.entry_profile.as_ref() {
4012 Some(profile) => self.resolve_profiled_entry(profile, &signal),
4013 None => resolve_unprofiled_entry(&signal),
4014 };
4015 match resolved {
4016 Ok(Some(resolved)) => {
4017 let side_for_audit = resolved.side;
4018 let level_resolution = resolved.level_resolution.clone();
4019 let resolved_target_prices: Vec<f64> =
4020 resolved.targets.iter().map(|target| target.price).collect();
4021 match self.finalize_resolved_entry(
4022 resolved,
4023 future_executor.balance(),
4024 quote.ts,
4025 Some(conversion_quotes),
4026 Some(sizing_reference_price),
4027 ) {
4028 Ok(finalized) => {
4029 let (trade_id, protective_stop) = match &finalized.action {
4030 Action::Open {
4031 trade_id, stoploss, ..
4032 } => (trade_id.clone(), *stoploss),
4033 _ => unreachable!("finalized entry must be an open action"),
4034 };
4035 let mut levels_crossed_at_fill = Vec::new();
4038 if let Some(stop) = protective_stop {
4039 let crossed = match side_for_audit {
4040 Side::Buy => execution.price <= stop,
4041 Side::Sell => execution.price >= stop,
4042 };
4043 if crossed {
4044 levels_crossed_at_fill.push("stop".to_owned());
4045 }
4046 }
4047 for (offset, target_price) in
4048 resolved_target_prices.iter().enumerate()
4049 {
4050 let crossed = match side_for_audit {
4051 Side::Buy => execution.price >= *target_price,
4052 Side::Sell => execution.price <= *target_price,
4053 };
4054 if crossed {
4055 levels_crossed_at_fill
4056 .push(format!("target{}", offset + 1));
4057 }
4058 }
4059 action.entry_profile_resolution_audit =
4060 Some(EntryProfileResolutionAudit {
4061 action_id: action.action_id.clone(),
4062 trade_id: trade_id.clone(),
4063 entry_class,
4064 selection_source: action
4065 .entry_profile_selection_source
4066 .unwrap_or(EntryProfileSelectionSource::Unprofiled),
4067 selected_profile_name: action.selected_profile_name.clone(),
4068 resolution_stage: EntryResolutionStage::MarketExecution,
4069 original_signal_price,
4070 level_reference_price: Some(execution.price),
4071 level_resolution: Some(finalized.level_resolution.clone()),
4072 target_resolution: Some(
4073 finalized.target_resolution.clone(),
4074 ),
4075 configured_weights: finalized.configured_weights.clone(),
4076 allocated_target_steps: finalized
4077 .allocated_target_steps
4078 .clone(),
4079 remainder_steps: finalized.remainder_steps,
4080 outcome: crate::ledger::ActionDispositionStatus::Applied,
4081 rejection_stage: None,
4082 reason: None,
4083 });
4084 action.market_entry_sizing_audit = Some(MarketEntrySizingAudit {
4085 action_id: action.action_id.clone(),
4086 trade_id,
4087 configured_basis,
4088 applied_basis,
4089 fallback_to_fill,
4090 original_signal_price,
4091 sizing_reference_price,
4092 execution_price: execution.price,
4093 protective_stop,
4094 requested_account_risk: finalized.requested_account_risk,
4095 native_loss_per_lot: finalized.native_loss_per_lot,
4096 account_loss_per_lot: finalized.account_loss_per_lot,
4097 final_lot: finalized.final_lot,
4098 levels_crossed_at_fill,
4099 });
4100 action.action = finalized.action;
4101 }
4102 Err(error) => {
4103 self.record_entry_resolution_rejection(
4104 action.action_id.clone(),
4105 &signal,
4106 Some(&SelectedEntryProfile {
4107 profile: action.entry_profile.clone(),
4108 source: action
4109 .entry_profile_selection_source
4110 .unwrap_or(EntryProfileSelectionSource::Unprofiled),
4111 entry_class: entry_class.clone(),
4112 profile_name: action.selected_profile_name.clone(),
4113 }),
4114 EntryResolutionStage::MarketExecution,
4115 Some(execution.price),
4116 Some(level_resolution),
4117 "sizing",
4118 error.clone(),
4119 );
4120 let mut disposition =
4121 ActionDisposition::rejected(action.action_id, error);
4122 disposition.action_kind = Some(action.action_kind);
4123 disposition.signal_ts = Some(action.signal_ts);
4124 disposition.effective_ts = Some(action.effective_ts);
4125 self.record_disposition(lifecycle, disposition);
4126 continue;
4127 }
4128 }
4129 }
4130 Ok(None) => {
4131 let mut disposition =
4132 ActionDisposition::skipped(action.action_id, "not_an_entry");
4133 disposition.action_kind = Some(action.action_kind);
4134 disposition.signal_ts = Some(action.signal_ts);
4135 disposition.effective_ts = Some(action.effective_ts);
4136 self.record_disposition(lifecycle, disposition);
4137 continue;
4138 }
4139 Err(error) => {
4140 let reason = error.to_string();
4141 self.record_entry_resolution_rejection(
4142 action.action_id.clone(),
4143 &signal,
4144 Some(&SelectedEntryProfile {
4145 profile: action.entry_profile.clone(),
4146 source: action
4147 .entry_profile_selection_source
4148 .unwrap_or(EntryProfileSelectionSource::Unprofiled),
4149 entry_class: entry_class.clone(),
4150 profile_name: action.selected_profile_name.clone(),
4151 }),
4152 EntryResolutionStage::MarketExecution,
4153 Some(execution.price),
4154 None,
4155 "profile_resolution",
4156 reason.clone(),
4157 );
4158 let mut disposition = ActionDisposition::rejected(action.action_id, reason);
4159 disposition.action_kind = Some(action.action_kind);
4160 disposition.signal_ts = Some(action.signal_ts);
4161 disposition.effective_ts = Some(action.effective_ts);
4162 self.record_disposition(lifecycle, disposition);
4163 continue;
4164 }
4165 }
4166 }
4167 let committed = self.apply_future_action(
4168 action.action_id,
4169 action.action_kind,
4170 action.action,
4171 action.execution,
4172 action.signal_ts,
4173 action.effective_ts,
4174 quote,
4175 lifecycle,
4176 future_executor,
4177 portfolio,
4178 pricer,
4179 conversion_quotes,
4180 );
4181 if committed {
4182 if let Some(audit) = action.market_entry_sizing_audit {
4183 self.market_entry_sizing.push(audit);
4184 }
4185 if let Some(audit) = action.entry_profile_resolution_audit {
4186 self.entry_profile_resolutions.push(audit);
4187 }
4188 } else if let Some(mut audit) = action.entry_profile_resolution_audit {
4189 if let Some(disposition) = lifecycle
4190 .as_slice()
4191 .iter()
4192 .rev()
4193 .find(|disposition| disposition.action_id == audit.action_id)
4194 {
4195 audit.outcome = disposition.status;
4196 audit.rejection_stage = Some("engine_or_accounting".into());
4197 audit.reason = disposition.reason.clone();
4198 }
4199 self.entry_profile_resolutions.push(audit);
4200 }
4201 }
4202 *queued = remaining;
4203 }
4204
4205 #[allow(clippy::too_many_arguments)]
4206 fn record_entry_resolution_rejection(
4207 &mut self,
4208 action_id: String,
4209 signal: &RawSignal,
4210 selected: Option<&SelectedEntryProfile>,
4211 stage: EntryResolutionStage,
4212 level_reference_price: Option<f64>,
4213 level_resolution: Option<qs_core::EntryLevelResolution>,
4214 rejection_stage: &str,
4215 reason: String,
4216 ) {
4217 let (trade_id, original_signal_price, entry_class) = match signal {
4218 RawSignal::Entry {
4219 trade_id,
4220 price,
4221 entry_class,
4222 ..
4223 } => (trade_id.clone(), *price, entry_class.clone()),
4224 _ => (None, None, None),
4225 };
4226 let selection_source = selected.map_or_else(
4227 || {
4228 if entry_class.is_some() {
4229 EntryProfileSelectionSource::Mapped
4230 } else {
4231 EntryProfileSelectionSource::Unprofiled
4232 }
4233 },
4234 |selection| selection.source,
4235 );
4236 self.entry_profile_resolutions
4237 .push(EntryProfileResolutionAudit {
4238 action_id,
4239 trade_id,
4240 entry_class,
4241 selection_source,
4242 selected_profile_name: selected
4243 .and_then(|selection| selection.profile_name.clone()),
4244 resolution_stage: stage,
4245 original_signal_price,
4246 level_reference_price,
4247 level_resolution,
4248 target_resolution: None,
4249 configured_weights: Vec::new(),
4250 allocated_target_steps: Vec::new(),
4251 remainder_steps: 0,
4252 outcome: ActionDispositionStatus::Rejected,
4253 rejection_stage: Some(rejection_stage.into()),
4254 reason: Some(reason),
4255 });
4256 }
4257
4258 fn record_disposition(
4259 &mut self,
4260 lifecycle: &mut LifecycleLedger,
4261 disposition: ActionDisposition,
4262 ) {
4263 if lifecycle.record(disposition.clone()).is_ok() {
4264 self.committed_feedback_events
4265 .push(StrategyFeedbackEvent::Disposition(disposition));
4266 }
4267 }
4268
4269 fn record_committed_effects(&mut self, effects: Vec<FutureEffect>, action_id: Option<String>) {
4270 for effect in effects {
4271 self.committed_feedback_events
4272 .push(StrategyFeedbackEvent::Effect {
4273 action_id: action_id.clone(),
4274 effect: effect.clone(),
4275 });
4276 self.committed_feedback.push(effect);
4277 }
4278 }
4279
4280 #[allow(clippy::too_many_arguments)]
4281 fn apply_future_action(
4282 &mut self,
4283 action_id: String,
4284 action_kind: String,
4285 mut action: Action,
4286 execution: Option<ExecutionFill>,
4287 signal_ts: NaiveDateTime,
4288 effective_ts: NaiveDateTime,
4289 quote: &PriceQuote,
4290 lifecycle: &mut LifecycleLedger,
4291 future_executor: &mut FutureExecutor,
4292 portfolio: &mut PortfolioRecorder,
4293 pricer: &ExecutionPricer,
4294 conversion_quotes: &ConversionQuoteBook,
4295 ) -> bool {
4296 if let Action::ScaleIn {
4297 position_id, size, ..
4298 } = &action
4299 {
4300 let valid = self
4301 .engine
4302 .get_position(position_id)
4303 .and_then(|position| explicit_instrument_spec(&self.config, &position.data.symbol))
4304 .is_none_or(|spec| {
4305 size.to_string().parse::<Decimal>().is_ok_and(|quantity| {
4306 quantity >= spec.quantity.minimum.get()
4307 && spec
4308 .quantity
4309 .maximum
4310 .is_none_or(|maximum| quantity <= maximum.get())
4311 && spec.quantity.grid.contains(quantity).unwrap_or(false)
4312 })
4313 });
4314 if !valid {
4315 let mut disposition =
4316 ActionDisposition::rejected(action_id, "invalid_instrument_quantity");
4317 disposition.action_kind = Some(action_kind);
4318 disposition.signal_ts = Some(signal_ts);
4319 disposition.effective_ts = Some(effective_ts);
4320 self.record_disposition(lifecycle, disposition);
4321 return false;
4322 }
4323 }
4324 if let Action::Open {
4325 trade_id: Some(trade_id),
4326 ..
4327 } = &action
4328 && self.engine.manager.id_by_trade_id(trade_id).is_some()
4329 {
4330 let mut disposition = ActionDisposition::rejected(action_id, "duplicate_trade_id");
4331 disposition.action_kind = Some(action_kind);
4332 disposition.signal_ts = Some(signal_ts);
4333 disposition.effective_ts = Some(effective_ts);
4334 self.record_disposition(lifecycle, disposition);
4335 return false;
4336 }
4337 if let Action::ScaleIn { position_id, .. } = &action
4338 && future_executor.has_close(position_id)
4339 {
4340 let mut disposition =
4341 ActionDisposition::rejected(action_id, "scale_in_after_close_not_supported");
4342 disposition.action_kind = Some(action_kind);
4343 disposition.signal_ts = Some(signal_ts);
4344 disposition.effective_ts = Some(effective_ts);
4345 disposition.position_ids.push(position_id.clone());
4346 self.record_disposition(lifecycle, disposition);
4347 return false;
4348 }
4349
4350 let execution = match self.prepare_future_action(&mut action, execution, quote, pricer) {
4351 Ok(execution) => execution,
4352 Err(reason) => {
4353 let mut disposition = ActionDisposition::rejected(action_id, reason);
4354 disposition.action_kind = Some(action_kind);
4355 disposition.signal_ts = Some(signal_ts);
4356 disposition.effective_ts = Some(effective_ts);
4357 self.record_disposition(lifecycle, disposition);
4358 return false;
4359 }
4360 };
4361
4362 let engine_transaction = match execution {
4363 Some(execution) => self
4364 .engine
4365 .begin_priced_future_action(action, quote, execution),
4366 None => self.engine.begin_future_action(action, effective_ts),
4367 };
4368 let engine_transaction = match engine_transaction {
4369 Ok(transaction) => transaction,
4370 Err(error) => {
4371 let closed_state = matches!(
4374 error,
4375 FutureApplyError::Core(qs_core::CoreError::InvalidState { .. })
4376 ) && action_kind != "entry";
4377 let mut disposition = if closed_state {
4378 ActionDisposition::skipped(action_id, "position_closed")
4379 } else {
4380 ActionDisposition::rejected(action_id, error.to_string())
4381 };
4382 disposition.action_kind = Some(action_kind);
4383 disposition.signal_ts = Some(signal_ts);
4384 disposition.effective_ts = Some(effective_ts);
4385 self.record_disposition(lifecycle, disposition);
4386 return false;
4387 }
4388 };
4389
4390 let committed_effects = engine_transaction.effects().to_vec();
4391 let mut affected = if FutureExecutor::requires_processing(engine_transaction.effects()) {
4392 match future_executor.process_future_effects_with_currency(
4393 engine_transaction.effects(),
4394 &self.engine,
4395 quote,
4396 Some(&action_id),
4397 Some(signal_ts),
4398 effective_ts,
4399 portfolio,
4400 Some(conversion_quotes),
4401 ) {
4402 Ok(affected) => affected,
4403 Err(error) => {
4404 engine_transaction.rollback(&mut self.engine);
4405 let mut disposition = ActionDisposition::failed(action_id, error.to_string());
4406 disposition.action_kind = Some(action_kind);
4407 disposition.signal_ts = Some(signal_ts);
4408 disposition.effective_ts = Some(effective_ts);
4409 self.record_disposition(lifecycle, disposition);
4410 return false;
4411 }
4412 }
4413 } else {
4414 Vec::new()
4415 };
4416 for future_effect in engine_transaction.effects() {
4417 match future_effect.effect() {
4418 Effect::OrderPlaced { id } | Effect::OrderCancelled { id } => {
4419 affected.push(id.clone());
4420 }
4421 _ => {}
4422 }
4423 }
4424 affected.sort();
4425 affected.dedup();
4426 let _ = engine_transaction.commit();
4427 self.record_committed_effects(committed_effects, Some(action_id.clone()));
4428
4429 let mut disposition = ActionDisposition::applied(action_id);
4430 disposition.action_kind = Some(action_kind);
4431 disposition.signal_ts = Some(signal_ts);
4432 disposition.effective_ts = Some(effective_ts);
4433 disposition.position_ids = affected;
4434 self.record_disposition(lifecycle, disposition);
4435 true
4436 }
4437
4438 fn prepare_future_action(
4439 &self,
4440 action: &mut Action,
4441 prepriced: Option<ExecutionFill>,
4442 quote: &PriceQuote,
4443 pricer: &ExecutionPricer,
4444 ) -> Result<Option<ExecutionFill>, String> {
4445 let mut execution = prepriced;
4446 match action {
4447 Action::Open {
4448 symbol,
4449 side,
4450 order_type,
4451 price,
4452 size,
4453 ..
4454 } => {
4455 if !valid_accounting_size(*size) {
4456 return Err(format!(
4457 "position size must be finite and greater than the accounting tolerance, got {size}"
4458 ));
4459 }
4460 if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
4461 return Err(format!(
4462 "supplied entry price must be finite and positive, got {price:?}"
4463 ));
4464 }
4465 if *order_type == OrderType::Market {
4466 let priced = match execution {
4467 Some(priced) => priced,
4468 None => pricer
4469 .market_entry(*side, quote, self.pip_size(symbol))
4470 .map_err(|error| error.to_string())?,
4471 };
4472 *price = Some(priced.price);
4473 execution = Some(priced);
4474 } else {
4475 if price.is_none() {
4476 return Err("pending entry requires a requested price".to_owned());
4477 }
4478 execution = None;
4479 }
4480 }
4481 Action::ScaleIn {
4482 position_id,
4483 price,
4484 size,
4485 ..
4486 } => {
4487 if !valid_accounting_size(*size) {
4488 return Err(format!(
4489 "scale-in size must be finite and greater than the accounting tolerance, got {size}"
4490 ));
4491 }
4492 if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
4493 return Err(format!(
4494 "supplied scale-in price must be finite and positive, got {price:?}"
4495 ));
4496 }
4497 let side = self
4498 .engine
4499 .get_position(position_id)
4500 .map(|position| position.data.side)
4501 .ok_or_else(|| format!("position not found: {position_id}"))?;
4502 let priced = match execution {
4503 Some(priced) => priced,
4504 None => pricer
4505 .market_entry(side, quote, self.pip_size("e.symbol))
4506 .map_err(|error| error.to_string())?,
4507 };
4508 *price = Some(priced.price);
4509 execution = Some(priced);
4510 }
4511 Action::ClosePosition { position_id } | Action::ClosePartial { position_id, .. } => {
4512 let position = self
4513 .engine
4514 .get_position(position_id)
4515 .ok_or_else(|| format!("position not found: {position_id}"))?;
4516 if position.data.symbol != quote.symbol {
4517 return Err(format!(
4518 "position symbol {} does not match quote symbol {}",
4519 position.data.symbol, quote.symbol
4520 ));
4521 }
4522 execution = Some(
4523 pricer
4524 .market_exit(
4525 position.data.side,
4526 quote,
4527 self.pip_size(&position.data.symbol),
4528 )
4529 .map_err(|error| error.to_string())?,
4530 );
4531 }
4532 _ => execution = None,
4533 }
4534 Ok(execution)
4535 }
4536
4537 fn prepare_triggering_pending(
4538 &self,
4539 quote: &PriceQuote,
4540 side: Option<Side>,
4541 pricer: &ExecutionPricer,
4542 ) -> (Vec<PreparedPendingFill>, Vec<(String, String)>) {
4543 let ids = self
4544 .engine
4545 .manager
4546 .pending_ids_by_symbol_sorted("e.symbol);
4547 let mut prepared = Vec::new();
4548 let mut failures = Vec::new();
4549 for id in ids {
4550 let Some(position) = self.engine.get_position(&id) else {
4551 continue;
4552 };
4553 if side.is_some_and(|side| position.data.side != side) {
4554 continue;
4555 }
4556 let Some(purpose) = position.pending_fill_purpose(quote, self.config.fill_model) else {
4557 continue;
4558 };
4559 let execution = match pricer.price(
4560 purpose,
4561 position.data.side,
4562 quote,
4563 position.data.pending_price,
4564 self.pip_size("e.symbol),
4565 ) {
4566 Ok(fill) => fill,
4567 Err(error) => {
4568 failures.push((id, error.to_string()));
4569 continue;
4570 }
4571 };
4572
4573 let size = position.data.size;
4574 if !valid_accounting_size(size) {
4575 failures.push((
4576 id,
4577 format!(
4578 "pending size must be finite and greater than the accounting tolerance, got {size}"
4579 ),
4580 ));
4581 continue;
4582 }
4583
4584 prepared.push(PreparedPendingFill {
4585 position_id: id,
4586 execution,
4587 size,
4588 });
4589 }
4590 (prepared, failures)
4591 }
4592
4593 fn pip_size(&self, symbol: &str) -> f64 {
4594 self.config
4595 .symbol_specs
4596 .get(symbol)
4597 .map(|spec| 10_f64.powi(-(spec.pip_position as i32)))
4598 .unwrap_or(0.0001)
4599 }
4600
4601 fn action_symbol(&self, action: &Action) -> Option<String> {
4602 match action {
4603 Action::Open { symbol, .. } => Some(symbol.clone()),
4604 Action::ClosePosition { position_id }
4605 | Action::ClosePartial { position_id, .. }
4606 | Action::ModifyStoploss { position_id, .. }
4607 | Action::MoveStoplossToEntry { position_id }
4608 | Action::AddTarget { position_id, .. }
4609 | Action::RemoveTarget { position_id, .. }
4610 | Action::ModifyTarget { position_id, .. }
4611 | Action::AddRule { position_id, .. }
4612 | Action::RemoveRule { position_id, .. }
4613 | Action::ScaleIn { position_id, .. }
4614 | Action::CancelPending { position_id } => self
4615 .engine
4616 .get_position(position_id)
4617 .map(|position| position.data.symbol.clone()),
4618 Action::CloseAllOf { symbol } | Action::ModifyAllStoploss { symbol, .. } => {
4619 Some(symbol.clone())
4620 }
4621 _ => None,
4622 }
4623 }
4624
4625 fn resolve_future_actions(&self, signal: &RawSignal) -> Vec<Action> {
4626 match signal {
4627 RawSignal::CloseAllOf { symbol, .. } => self
4628 .engine
4629 .manager
4630 .open_ids_by_symbol_sorted(symbol)
4631 .into_iter()
4632 .map(|position_id| Action::ClosePosition { position_id })
4633 .collect(),
4634 RawSignal::CloseAll { .. } => self
4635 .engine
4636 .manager
4637 .ids_by_status_sorted(PositionStatus::Open)
4638 .into_iter()
4639 .map(|position_id| Action::ClosePosition { position_id })
4640 .collect(),
4641 RawSignal::CloseAllInGroup { group_id, .. } => {
4642 let mut ids = self.engine.manager.open_ids_by_group(group_id);
4643 ids.sort();
4644 ids.into_iter()
4645 .map(|position_id| Action::ClosePosition { position_id })
4646 .collect()
4647 }
4648 RawSignal::CancelAllPending { .. } => self
4649 .engine
4650 .manager
4651 .ids_by_status_sorted(PositionStatus::Pending)
4652 .into_iter()
4653 .map(|position_id| Action::CancelPending { position_id })
4654 .collect(),
4655 _ => resolve_signal(signal, &self.engine),
4656 }
4657 }
4658}
4659
4660fn map_strategy_driver_error<FeedError, StrategyError>(
4661 error: StrategyDriverError<StrategyError>,
4662) -> StrategyReplayError<FeedError, StrategyError> {
4663 match error {
4664 StrategyDriverError::Series(error) => StrategyReplayError::Series(error),
4665 StrategyDriverError::SeriesView(error) => StrategyReplayError::SeriesView(error),
4666 StrategyDriverError::Analysis(error) => StrategyReplayError::Analysis(error),
4667 StrategyDriverError::Strategy(error) => StrategyReplayError::Strategy(error),
4668 StrategyDriverError::Runtime(error) => StrategyReplayError::Runtime(error),
4669 StrategyDriverError::WarmupSignals { timestamp } => {
4670 StrategyReplayError::WarmupSignals { timestamp }
4671 }
4672 StrategyDriverError::InvalidGeneratedSignal {
4673 signal_index,
4674 reason,
4675 } => StrategyReplayError::InvalidGeneratedSignal {
4676 signal_index,
4677 reason,
4678 },
4679 StrategyDriverError::TickExecutionRequired { symbol, timestamp } => {
4680 StrategyReplayError::TickExecutionRequired { symbol, timestamp }
4681 }
4682 }
4683}
4684
4685#[allow(clippy::too_many_arguments)]
4686fn observe_future_equity(
4687 portfolio: &mut PortfolioRecorder,
4688 future_executor: &FutureExecutor,
4689 ts: NaiveDateTime,
4690 conversion_quotes: &ConversionQuoteBook,
4691 kind: EquityObservationKind,
4692 collector: &mut MtmCurveCollector,
4693 last_candidate: &mut Option<EquityPoint>,
4694 suppress_unchanged_post_output: bool,
4695) {
4696 portfolio.set_realized_pnl(future_executor.realized_pnl());
4697 let mut point = portfolio.observe_with_currency(
4698 ts,
4699 future_executor.open_snapshots(),
4700 Some(conversion_quotes),
4701 );
4702 point.observation_kind = Some(kind.as_str().to_owned());
4703 if suppress_unchanged_post_output
4704 && last_candidate
4705 .as_ref()
4706 .is_some_and(|previous| same_equity_values(previous, &point))
4707 {
4708 return;
4709 }
4710 collector.observe(point.clone());
4711 *last_candidate = Some(point);
4712}
4713
4714fn same_equity_values(left: &EquityPoint, right: &EquityPoint) -> bool {
4715 let mut left = left.clone();
4716 let mut right = right.clone();
4717 left.observation_kind = None;
4718 left.observation_sequence = None;
4719 right.observation_kind = None;
4720 right.observation_sequence = None;
4721 left == right
4722}
4723
4724fn pending_fill_position_id(effect: &FutureEffect) -> Option<&str> {
4725 match effect.effect() {
4726 Effect::PositionOpened { id } => Some(id),
4727 _ => None,
4728 }
4729}
4730
4731fn valid_accounting_size(size: f64) -> bool {
4732 size.is_finite() && size > position_size_tolerance(size)
4733}
4734
4735fn explicit_instrument_spec<'a>(
4736 config: &'a BacktestConfig,
4737 symbol: &str,
4738) -> Option<&'a InstrumentSpec> {
4739 config
4740 .instrument_manifest
4741 .as_ref()?
4742 .instruments
4743 .get(symbol)
4744 .map(|artifact| &artifact.spec)
4745}
4746
4747fn decimal_to_f64(value: Decimal, field: &str) -> Result<f64, String> {
4748 let value = value
4749 .to_string()
4750 .parse::<f64>()
4751 .map_err(|error| format!("invalid {field}: {error}"))?;
4752 if value.is_finite() {
4753 Ok(value)
4754 } else {
4755 Err(format!("{field} must be finite"))
4756 }
4757}
4758
4759fn instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
4760 decimal_to_f64(
4761 spec.economics.contract_multiplier.get(),
4762 "instrument contract multiplier",
4763 )
4764 .and_then(|value| {
4765 if value > 0.0 {
4766 Ok(value)
4767 } else {
4768 Err("instrument contract multiplier must be positive".into())
4769 }
4770 })
4771}
4772
4773fn supported_instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
4774 if spec.status != ListingStatus::Trading {
4775 return Err(format!(
4776 "instrument {} is not in trading status",
4777 spec.instrument
4778 ));
4779 }
4780 if spec.economics.quantity_unit != QuantityUnit::StandardLot {
4781 return Err(format!(
4782 "unsupported quantity unit for instrument {}: {:?}",
4783 spec.instrument, spec.economics.quantity_unit
4784 ));
4785 }
4786 let model = spec.economics.pnl_model.as_str();
4787 if model != EconomicsModelId::FX_QUOTE_LINEAR_V1
4788 && model != EconomicsModelId::CFD_QUOTE_LINEAR_V1
4789 {
4790 return Err(format!(
4791 "unsupported P&L model for instrument {}: {model}",
4792 spec.instrument
4793 ));
4794 }
4795 instrument_multiplier(spec)
4796}
4797
4798fn validate_instrument_manifest(config: &BacktestConfig) -> Result<(), String> {
4799 let Some(manifest) = &config.instrument_manifest else {
4800 return Ok(());
4801 };
4802 for (symbol, artifact) in &manifest.instruments {
4803 if symbol.is_empty() {
4804 return Err("instrument manifest symbol must not be empty".into());
4805 }
4806 artifact
4807 .spec
4808 .validate()
4809 .map_err(|error| format!("invalid instrument spec for {symbol}: {error}"))?;
4810 if artifact.resolved.instrument != artifact.spec.instrument {
4811 return Err(format!(
4812 "resolved instrument and spec identity differ for {symbol}"
4813 ));
4814 }
4815 if artifact.resolved.spec_revision != artifact.spec.revision {
4816 return Err(format!(
4817 "resolved specification revision does not match the embedded spec for {symbol}"
4818 ));
4819 }
4820 supported_instrument_multiplier(&artifact.spec)?;
4821 }
4822 for binding in &manifest.stored_series {
4823 let known = manifest
4824 .instruments
4825 .values()
4826 .any(|artifact| artifact.resolved == binding.instrument);
4827 if !known {
4828 return Err(format!(
4829 "stored series {}:{} references an instrument outside the manifest",
4830 binding.source_partition, binding.source_symbol
4831 ));
4832 }
4833 let artifact = manifest
4834 .instruments
4835 .values()
4836 .find(|artifact| artifact.resolved == binding.instrument)
4837 .expect("known binding reference has an instrument artifact");
4838 if binding.effective != artifact.spec.effective {
4839 return Err(format!(
4840 "stored series {}:{} effective interval differs from its instrument spec",
4841 binding.source_partition, binding.source_symbol
4842 ));
4843 }
4844 }
4845 Ok(())
4846}
4847
4848fn effective_contract_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
4849 let mut contract_sizes = config.contract_sizes.clone();
4850 if let Some(manifest) = &config.instrument_manifest {
4851 for (symbol, artifact) in &manifest.instruments {
4852 if let Ok(multiplier) = instrument_multiplier(&artifact.spec) {
4853 contract_sizes.insert(symbol.clone(), multiplier);
4854 }
4855 }
4856 }
4857 contract_sizes
4858}
4859
4860fn validate_replay_costs(
4862 config: &BacktestConfig,
4863 future: Option<&FutureQuoteConfig>,
4864) -> Result<(), String> {
4865 let account_currency = future
4866 .and_then(|future| future.currency_plan.as_ref())
4867 .map(|plan| plan.account_currency().to_owned());
4868 for (symbol, costs) in &config.costs {
4869 if symbol.is_empty() {
4870 return Err("cost symbol must not be empty".into());
4871 }
4872 costs
4873 .validate()
4874 .map_err(|error| format!("costs for {symbol} are invalid: {error}"))?;
4875 if let Some(account_currency) = account_currency.as_deref() {
4876 costs
4877 .validate_against_account_currency(account_currency)
4878 .map_err(|error| format!("costs for {symbol} are invalid: {error}"))?;
4879 }
4880 if costs.requires_point_size() && !config.symbol_specs.contains_key(symbol) {
4881 return Err(format!(
4882 "point-denominated swap for {symbol} requires a symbol specification for its digit count"
4883 ));
4884 }
4885 }
4886 Ok(())
4887}
4888
4889fn effective_point_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
4891 config
4892 .symbol_specs
4893 .iter()
4894 .map(|(symbol, spec)| (symbol.clone(), 10f64.powi(-i32::from(spec.digits))))
4895 .collect()
4896}
4897
4898fn is_monetary_sizing(policy: &SizingPolicy) -> bool {
4899 matches!(
4900 policy,
4901 SizingPolicy::FixedRiskAmount { .. } | SizingPolicy::BalanceRiskPercent { .. }
4902 )
4903}
4904
4905fn accept_legacy_quote(
4906 quote: &PriceQuote,
4907 last_quote_ts: &mut BTreeMap<String, NaiveDateTime>,
4908) -> bool {
4909 if ExecutionPricer::validate_quote(quote).is_err()
4910 || last_quote_ts
4911 .get("e.symbol)
4912 .is_some_and(|last| *last > quote.ts)
4913 {
4914 return false;
4915 }
4916 last_quote_ts.insert(quote.symbol.clone(), quote.ts);
4917 true
4918}
4919
4920pub const MAX_RUN_TAGS: usize = 32;
4922pub const MAX_RUN_TAG_BYTES: usize = 64;
4924
4925fn validate_run_tags(tags: &BTreeMap<String, String>) -> Result<(), String> {
4926 if tags.len() > MAX_RUN_TAGS {
4927 return Err(format!(
4928 "run tags must not exceed {MAX_RUN_TAGS} entries, got {}",
4929 tags.len()
4930 ));
4931 }
4932 for (key, value) in tags {
4933 if key.is_empty() || key.len() > MAX_RUN_TAG_BYTES {
4934 return Err(format!(
4935 "run tag key must be 1 to {MAX_RUN_TAG_BYTES} bytes, got '{key}'"
4936 ));
4937 }
4938 if !key
4939 .chars()
4940 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
4941 {
4942 return Err(format!(
4943 "run tag key must be ASCII alphanumeric, underscore, or hyphen, got '{key}'"
4944 ));
4945 }
4946 if value.len() > MAX_RUN_TAG_BYTES {
4947 return Err(format!(
4948 "run tag value for '{key}' must not exceed {MAX_RUN_TAG_BYTES} bytes"
4949 ));
4950 }
4951 if value.chars().any(char::is_control) {
4952 return Err(format!(
4953 "run tag value for '{key}' must not contain control characters"
4954 ));
4955 }
4956 }
4957 Ok(())
4958}
4959
4960fn validate_replay_config(
4961 config: &BacktestConfig,
4962 future: Option<&FutureQuoteConfig>,
4963 raw_signals: &[RawSignal],
4964) -> Result<(), String> {
4965 if !config.initial_balance.is_finite() || config.initial_balance <= 0.0 {
4966 return Err(format!(
4967 "initial balance must be finite and positive, got {}",
4968 config.initial_balance
4969 ));
4970 }
4971 for (symbol, contract_size) in &config.contract_sizes {
4972 if symbol.is_empty() {
4973 return Err("contract-size symbol must not be empty".into());
4974 }
4975 if !contract_size.is_finite() || *contract_size <= 0.0 {
4976 return Err(format!(
4977 "contract size for {symbol} must be finite and positive, got {contract_size}"
4978 ));
4979 }
4980 }
4981
4982 validate_run_tags(&config.run_tags)?;
4983 validate_replay_costs(config, future)?;
4984 validate_instrument_manifest(config)?;
4985 for (symbol, spec) in &config.symbol_specs {
4986 validate_symbol_spec(symbol, spec)?;
4987 if explicit_instrument_spec(config, symbol).is_none() {
4988 resolve_legacy_economics(spec).map_err(|error| error.to_string())?;
4989 }
4990 }
4991
4992 let entry_symbols: Vec<&str> = raw_signals
4993 .iter()
4994 .filter_map(|signal| match signal {
4995 RawSignal::Entry { symbol, .. } => Some(symbol.as_str()),
4996 _ => None,
4997 })
4998 .collect();
4999 if !entry_symbols.is_empty() && config.sizing.is_none() {
5000 return Err("raw entry requires BacktestConfig.sizing".to_owned());
5001 }
5002 if let Some(policy) = &config.sizing {
5003 validate_sizing_policy(policy)?;
5004 for symbol in &entry_symbols {
5005 if !config.symbol_specs.contains_key(*symbol)
5006 && explicit_instrument_spec(config, symbol).is_none()
5007 {
5008 return Err(format!("missing instrument or symbol spec for {symbol}"));
5009 }
5010 }
5011 if !entry_symbols.is_empty() && is_monetary_sizing(policy) {
5012 let future = future.ok_or_else(|| {
5013 "monetary sizing requires FutureQuote execution and a currency plan".to_owned()
5014 })?;
5015 let plan = future
5016 .currency_plan
5017 .as_ref()
5018 .ok_or_else(|| "monetary sizing requires a FutureQuote currency plan".to_owned())?;
5019 for symbol in &entry_symbols {
5020 if plan.route_for_primary_symbol(symbol).is_none() {
5021 return Err(format!(
5022 "currency plan has no frozen route for primary symbol {symbol}"
5023 ));
5024 }
5025 }
5026 }
5027 }
5028
5029 if let Some(future) = future {
5030 if future.signal_latency_ms < 0 {
5031 return Err(format!(
5032 "signal latency must be non-negative, got {}",
5033 future.signal_latency_ms
5034 ));
5035 }
5036 let latency = Duration::milliseconds(future.signal_latency_ms);
5037 for signal in raw_signals {
5038 if signal.ts().checked_add_signed(latency).is_none() {
5039 return Err(format!(
5040 "signal latency overflows datetime for signal at {}",
5041 signal.ts()
5042 ));
5043 }
5044 }
5045 if !future.slippage_pips.is_finite() {
5046 return Err(format!(
5047 "slippage pips must be finite, got {}",
5048 future.slippage_pips
5049 ));
5050 }
5051 if future.stale_quote_after_ms.is_some_and(|value| value < 0) {
5052 return Err("stale quote threshold must be non-negative".into());
5053 }
5054 if !future.pnl_epsilon.is_finite() || future.pnl_epsilon < 0.0 {
5055 return Err(format!(
5056 "P&L epsilon must be finite and non-negative, got {}",
5057 future.pnl_epsilon
5058 ));
5059 }
5060 if future.conversion_stale_after_ms < 0 {
5061 return Err("conversion quote threshold must be non-negative".to_owned());
5062 }
5063 future
5064 .mtm_output
5065 .validate()
5066 .map_err(|error| error.to_string())?;
5067 }
5068 Ok(())
5069}
5070
5071fn validate_sizing_policy(policy: &SizingPolicy) -> Result<(), String> {
5072 let (name, value) = match policy {
5073 SizingPolicy::FixedLot { lots } => ("fixed lots", *lots),
5074 SizingPolicy::FixedRiskAmount { amount } => ("fixed risk amount", *amount),
5075 SizingPolicy::BalanceRiskPercent { percent } => ("balance risk percent", *percent),
5076 };
5077 if value.is_finite() && value > 0.0 {
5078 Ok(())
5079 } else {
5080 Err(format!("{name} must be finite and positive, got {value}"))
5081 }
5082}
5083
5084fn validate_symbol_spec(symbol: &str, spec: &qs_symbols::SymbolSpec) -> Result<(), String> {
5085 if symbol.is_empty() || spec.canonical.is_empty() {
5086 return Err("symbol spec names must not be empty".into());
5087 }
5088 if spec.digits > 18 || spec.pip_position > spec.digits {
5089 return Err(format!(
5090 "invalid price precision for {symbol}: digits={}, pip_position={}",
5091 spec.digits, spec.pip_position
5092 ));
5093 }
5094 if spec.lot_base_units <= 0
5095 || spec.lot_step_units <= 0
5096 || spec.lot_min_steps <= 0
5097 || spec.lot_max_steps < 0
5098 || (spec.lot_max_steps > 0 && spec.lot_max_steps < spec.lot_min_steps)
5099 {
5100 return Err(format!("invalid lot metadata for {symbol}"));
5101 }
5102 let lot_step = spec.lot_step();
5103 let min_lot = spec.lot_min();
5104 let max_lot = spec.lot_max();
5105 if !lot_step.is_finite()
5106 || lot_step <= 0.0
5107 || !min_lot.is_finite()
5108 || min_lot <= 0.0
5109 || !max_lot.is_finite()
5110 {
5111 return Err(format!("invalid derived lot metadata for {symbol}"));
5112 }
5113 Ok(())
5114}
5115
5116fn rejected_legacy_result(config: &BacktestConfig) -> BacktestResult {
5117 BacktestResult::from_trade_log(
5118 if config.initial_balance.is_finite() {
5119 config.initial_balance
5120 } else {
5121 0.0
5122 },
5123 Vec::new(),
5124 )
5125}
5126
5127fn rejected_future_result(
5128 config: &BacktestConfig,
5129 future: &FutureQuoteConfig,
5130 evaluation_options: EvaluationOptions,
5131 error: String,
5132) -> BacktestResult {
5133 let execution_model = ExecutionModel::new(
5134 qs_core::types::ExecutionConvention::FutureQuoteV1,
5135 config.fill_model,
5136 if future.slippage_pips == 0.0 {
5137 SlippageModel::None
5138 } else {
5139 SlippageModel::FixedPips {
5140 pips: future.slippage_pips,
5141 }
5142 },
5143 );
5144 let mut lifecycle = LifecycleLedger::new();
5145 let _ = lifecycle.record(ActionDisposition::rejected(
5146 "configuration",
5147 format!("invalid_configuration: {error}"),
5148 ));
5149 let mut tags = BTreeMap::new();
5150 tags.insert("configuration_error".into(), error);
5151 insert_economic_support_metadata(&mut tags, config);
5152 let artifacts = FutureBacktestArtifacts {
5153 execution: ExecutionMetadata {
5154 execution_model,
5155 initial_balance: if config.initial_balance.is_finite() {
5156 config.initial_balance
5157 } else {
5158 0.0
5159 },
5160 account_currency: future
5161 .currency_plan
5162 .as_ref()
5163 .map(|plan| plan.account_currency().to_owned()),
5164 currency_plan: future.currency_plan.clone(),
5165 contract_sizes: effective_contract_sizes(config)
5166 .into_iter()
5167 .filter(|(symbol, size)| !symbol.is_empty() && size.is_finite() && *size > 0.0)
5168 .collect(),
5169 instrument_manifest: config.instrument_manifest.clone(),
5170 instrument_sizing: Vec::new(),
5171 market_entry_sizing_basis: future.market_entry_sizing_basis,
5172 market_entry_sizing: Vec::new(),
5173 stale_quote_after_millis: future.stale_quote_after_ms,
5174 pnl_epsilon: if future.pnl_epsilon.is_finite() && future.pnl_epsilon >= 0.0 {
5175 future.pnl_epsilon
5176 } else {
5177 crate::artifacts::DEFAULT_PNL_EPSILON
5178 },
5179 tags,
5180 ..ExecutionMetadata::default()
5181 },
5182 lifecycle,
5183 mtm_output_summary: MtmOutputSummary {
5184 policy: future.mtm_output,
5185 ..MtmOutputSummary::default()
5186 },
5187 ..FutureBacktestArtifacts::default()
5188 };
5189 BacktestResult::from_future_artifacts_with_options(artifacts, evaluation_options)
5190}
5191
5192fn insert_economic_support_metadata(tags: &mut BTreeMap<String, String>, config: &BacktestConfig) {
5193 let mut compatibility_specs = config.symbol_specs.iter().peekable();
5194 if compatibility_specs.peek().is_none() {
5195 return;
5196 }
5197 tags.insert(
5198 "economics.guard".into(),
5199 LEGACY_ECONOMIC_GUARD_ID.to_owned(),
5200 );
5201 for (symbol, spec) in compatibility_specs {
5202 let prefix = format!("economics.symbol.{symbol}");
5203 tags.insert(format!("{prefix}.category"), spec.category.clone());
5204 match resolve_legacy_economics(spec) {
5205 Ok(economics) => {
5206 tags.insert(format!("{prefix}.status"), "supported".into());
5207 tags.insert(format!("{prefix}.model"), economics.model.as_str().into());
5208 tags.insert(
5209 format!("{prefix}.contract_multiplier"),
5210 economics.contract_multiplier.to_string(),
5211 );
5212 }
5213 Err(error) => {
5214 tags.insert(format!("{prefix}.status"), "unsupported".into());
5215 tags.insert(format!("{prefix}.reason"), error.to_string());
5216 }
5217 }
5218 }
5219}
5220
5221fn queued_exposure_symbols(
5222 queued: &VecDeque<QueuedAction>,
5223 quotes: &BTreeMap<String, PriceQuote>,
5224 batch_ts: NaiveDateTime,
5225) -> BTreeSet<String> {
5226 queued
5227 .iter()
5228 .filter(|action| {
5229 action.effective_ts <= batch_ts
5230 && quotes.contains_key(&action.symbol)
5231 && is_exposure_increasing(&action.action)
5232 })
5233 .map(|action| action.symbol.clone())
5234 .collect()
5235}
5236
5237fn is_exposure_increasing(action: &Action) -> bool {
5238 matches!(
5239 action,
5240 Action::Open {
5241 order_type: OrderType::Market,
5242 ..
5243 } | Action::ScaleIn { .. }
5244 )
5245}
5246
5247fn is_fill_bearing(action: &Action) -> bool {
5248 matches!(
5249 action,
5250 Action::Open {
5251 order_type: OrderType::Market,
5252 ..
5253 } | Action::ClosePosition { .. }
5254 | Action::ClosePartial { .. }
5255 | Action::ScaleIn { .. }
5256 )
5257}
5258
5259fn raw_signal_kind(signal: &RawSignal) -> &'static str {
5260 match signal {
5261 RawSignal::Entry { .. } => "entry",
5262 RawSignal::Close { .. } => "close",
5263 RawSignal::ClosePartial { .. } => "close_partial",
5264 RawSignal::ModifyStoploss { .. } => "modify_stoploss",
5265 RawSignal::MoveStoplossToEntry { .. } => "move_stoploss_to_entry",
5266 RawSignal::AddTarget { .. } => "add_target",
5267 RawSignal::RemoveTarget { .. } => "remove_target",
5268 RawSignal::ModifyTarget { .. } => "modify_target",
5269 RawSignal::AddRule { .. } => "add_rule",
5270 RawSignal::RemoveRule { .. } => "remove_rule",
5271 RawSignal::ScaleIn { .. } => "scale_in",
5272 RawSignal::CancelPending { .. } => "cancel_pending",
5273 RawSignal::CloseAllOf { .. } => "close_all_of",
5274 RawSignal::CloseAll { .. } => "close_all",
5275 RawSignal::CancelAllPending { .. } => "cancel_all_pending",
5276 RawSignal::ModifyAllStoploss { .. } => "modify_all_stoploss",
5277 RawSignal::CloseAllInGroup { .. } => "close_all_in_group",
5278 RawSignal::ModifyAllStoplossInGroup { .. } => "modify_all_stoploss_in_group",
5279 }
5280}
5281
5282#[cfg(test)]
5285mod tests {
5286 use super::*;
5287 use crate::currency::{ConversionRoute, FxPair};
5288 use crate::data_feed::{EventMetadata, FeedEvent, MarketEvent, SeriesRoles, VecFeed};
5289 use crate::profile::{
5290 EntryGeometryPolicy, ManagementProfile, PositionRef, RawSignal, StoplossMode, TargetSource,
5291 };
5292 use chrono::NaiveDate;
5293 use qs_core::types::{CloseReason, FillPurpose, OrderType, Side, TargetSpec};
5294
5295 fn ts(h: u32, m: u32, s: u32) -> chrono::NaiveDateTime {
5296 NaiveDate::from_ymd_opt(2026, 1, 1)
5297 .unwrap()
5298 .and_hms_opt(h, m, s)
5299 .unwrap()
5300 }
5301
5302 fn tick(symbol: &str, bid: f64, ask: f64, time: chrono::NaiveDateTime) -> MarketEvent {
5303 MarketEvent::Tick {
5304 symbol: symbol.into(),
5305 ts: time,
5306 bid,
5307 ask,
5308 }
5309 }
5310
5311 #[test]
5312 fn scheduled_signal_preserves_an_explicit_opaque_action_id() {
5313 let scheduled = ScheduledSignal::new(
5314 7,
5315 ts(10, 0, 0),
5316 ts(10, 0, 1),
5317 RawSignal::CloseAll { ts: ts(10, 0, 0) },
5318 true,
5319 )
5320 .with_action_id("caller-command/opaque:7");
5321
5322 assert_eq!(scheduled.resolved_action_id(), "caller-command/opaque:7");
5323 }
5324
5325 #[test]
5326 fn scheduled_signal_keeps_the_compatible_generated_action_id() {
5327 let scheduled = ScheduledSignal::new(
5328 7,
5329 ts(10, 0, 0),
5330 ts(10, 0, 1),
5331 RawSignal::CloseAll { ts: ts(10, 0, 0) },
5332 false,
5333 );
5334
5335 assert_eq!(scheduled.resolved_action_id(), "signal:00000007");
5336 }
5337
5338 fn test_symbol_spec(symbol: &str) -> qs_symbols::SymbolSpec {
5339 qs_symbols::SymbolSpec {
5340 canonical: symbol.to_ascii_lowercase(),
5341 pip_position: 4,
5342 digits: 5,
5343 category: "forex".into(),
5344 lot_base_units: 100,
5345 lot_step_units: 1,
5346 lot_min_steps: 1,
5347 lot_max_steps: 0,
5348 }
5349 }
5350
5351 fn fixed_lot_config() -> BacktestConfig {
5352 BacktestConfig {
5353 sizing: Some(SizingPolicy::FixedLot { lots: 1.0 }),
5354 symbol_specs: ["EURUSD", "XAUUSD"]
5355 .into_iter()
5356 .map(|symbol| (symbol.to_owned(), test_symbol_spec(symbol)))
5357 .collect(),
5358 ..BacktestConfig::default()
5359 }
5360 }
5361
5362 fn identity_currency_plan(symbol: &str) -> RunCurrencyPlan {
5363 RunCurrencyPlan::new(
5364 "USD",
5365 [symbol.to_owned()].into_iter().collect(),
5366 Default::default(),
5367 [(symbol.to_owned(), "USD".to_owned())]
5368 .into_iter()
5369 .collect(),
5370 [(
5371 "USD".to_owned(),
5372 ConversionRoute::Identity {
5373 currency: "USD".to_owned(),
5374 },
5375 )]
5376 .into_iter()
5377 .collect(),
5378 Vec::new(),
5379 )
5380 .unwrap()
5381 }
5382
5383 struct ScriptedBatchFeed {
5384 batches: VecDeque<Result<Option<TimestampBatch>, &'static str>>,
5385 }
5386
5387 impl FallibleBatchFeed for ScriptedBatchFeed {
5388 type Error = &'static str;
5389
5390 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
5391 self.batches.pop_front().unwrap_or(Ok(None))
5392 }
5393 }
5394
5395 struct CountingBatchFeed {
5396 batches: VecDeque<TimestampBatch>,
5397 polls: std::rc::Rc<std::cell::Cell<usize>>,
5398 }
5399
5400 impl FallibleBatchFeed for CountingBatchFeed {
5401 type Error = Infallible;
5402
5403 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
5404 self.polls.set(self.polls.get() + 1);
5405 Ok(self.batches.pop_front())
5406 }
5407 }
5408
5409 fn primary_batch(event: MarketEvent) -> TimestampBatch {
5410 TimestampBatch {
5411 ts: event.ts(),
5412 events: vec![FeedEvent::new(
5413 event,
5414 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
5415 )],
5416 }
5417 }
5418
5419 fn market_entry(timestamp: NaiveDateTime, symbol: &str, order_type: OrderType) -> RawSignal {
5420 RawSignal::Entry {
5421 ts: timestamp,
5422 symbol: symbol.into(),
5423 side: Side::Buy,
5424 order_type,
5425 price: (order_type == OrderType::Limit).then_some(1.0),
5426 risk_multiplier: 1.0,
5427 stoploss: None,
5428 targets: Vec::new(),
5429 group: None,
5430 trade_id: Some(format!("{symbol}-blocker")),
5431 entry_class: None,
5432 }
5433 }
5434
5435 #[test]
5436 fn future_streaming_matches_materialized_and_stops_without_draining() {
5437 let events = vec![
5438 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5439 tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
5440 tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
5441 ];
5442 let signals = vec![
5443 market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market),
5444 RawSignal::CloseAll { ts: ts(10, 0, 1) },
5445 ];
5446 let config = BacktestConfig {
5447 close_on_finish: false,
5448 ..fixed_lot_config()
5449 };
5450 let mut materialized_feed = VecFeed::new(events.clone());
5451 let materialized = BacktestRunner::new_future(
5452 config.clone(),
5453 FutureQuoteConfig {
5454 mtm_output: MtmOutputPolicy::Full,
5455 ..FutureQuoteConfig::default()
5456 },
5457 )
5458 .run_raw_signals_future(&mut materialized_feed, signals.clone(), None);
5459
5460 let mut stream = ScriptedBatchFeed {
5461 batches: VecDeque::from([
5462 Ok(Some(primary_batch(events[0].clone()))),
5463 Ok(Some(primary_batch(events[1].clone()))),
5464 Ok(Some(primary_batch(events[2].clone()))),
5465 Err("must not drain"),
5466 ]),
5467 };
5468 let mut progress = Vec::new();
5469 let streamed = BacktestRunner::new_future(
5470 config,
5471 FutureQuoteConfig {
5472 mtm_output: MtmOutputPolicy::Full,
5473 ..FutureQuoteConfig::default()
5474 },
5475 )
5476 .run_raw_signals_future_streaming_controlled(
5477 &mut stream,
5478 Some(ts(10, 0, 2)),
5479 signals,
5480 None,
5481 || false,
5482 |update| progress.push(update),
5483 )
5484 .unwrap();
5485
5486 assert_eq!(
5487 serde_json::to_value(&streamed).unwrap(),
5488 serde_json::to_value(&materialized).unwrap()
5489 );
5490 assert_eq!(
5491 stream.batches.len(),
5492 2,
5493 "quiescence must leave the tail unread"
5494 );
5495 assert_eq!(progress.first().unwrap().total_events, 0);
5496 assert_eq!(progress.last().unwrap().processed_events, 2);
5497 assert_eq!(progress.last().unwrap().total_events, 2);
5498 assert_eq!(
5499 streamed.mtm_equity_curve.last().unwrap().ts,
5500 ts(10, 0, 1),
5501 "terminal observation must use the last processed primary timestamp"
5502 );
5503 assert_eq!(
5504 streamed
5505 .mtm_equity_curve
5506 .last()
5507 .unwrap()
5508 .observation_kind
5509 .as_deref(),
5510 Some(EquityObservationKind::QuiescentTermination.as_str())
5511 );
5512 assert_eq!(
5513 streamed
5514 .execution_metadata
5515 .as_ref()
5516 .unwrap()
5517 .tags
5518 .get("termination_reason")
5519 .map(String::as_str),
5520 Some("quiescent")
5521 );
5522 }
5523
5524 #[test]
5525 fn exact_time_close_waits_for_later_symbol_pending_fill() {
5526 let open_ts = ts(10, 0, 0);
5527 let execution_ts = ts(10, 0, 1);
5528 let events = vec![
5529 FeedEvent::new(
5530 tick("XAUUSD", 101.0, 101.0, open_ts),
5531 EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
5532 ),
5533 FeedEvent::new(
5534 tick("EURUSD", 1.1, 1.1, execution_ts),
5535 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
5536 ),
5537 FeedEvent::new(
5538 tick("XAUUSD", 100.0, 100.0, execution_ts),
5539 EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
5540 ),
5541 ];
5542 let signals = vec![
5543 RawSignal::Entry {
5544 ts: open_ts,
5545 symbol: "XAUUSD".into(),
5546 side: Side::Buy,
5547 order_type: OrderType::Limit,
5548 price: Some(100.0),
5549 risk_multiplier: 1.0,
5550 stoploss: None,
5551 targets: Vec::new(),
5552 group: None,
5553 trade_id: Some("later-pending".into()),
5554 entry_class: None,
5555 },
5556 RawSignal::Close {
5557 ts: execution_ts,
5558 position: PositionRef::ByTradeId {
5559 trade_id: "later-pending".into(),
5560 },
5561 },
5562 ];
5563 let mut feed = VecFeed::from_feed_events(events);
5564 let result = BacktestRunner::new_future(
5565 BacktestConfig {
5566 close_on_finish: false,
5567 ..fixed_lot_config()
5568 },
5569 FutureQuoteConfig::default(),
5570 )
5571 .run_raw_signals_future(&mut feed, signals, None);
5572
5573 assert_eq!(
5574 result
5575 .recorded_fills
5576 .iter()
5577 .map(|fill| fill.fill.purpose)
5578 .collect::<Vec<_>>(),
5579 vec![FillPurpose::LimitEntry, FillPurpose::MarketExit]
5580 );
5581 assert_eq!(result.close_events.len(), 1);
5582 assert_eq!(result.close_events[0].reason, CloseReason::Manual);
5583 assert!(result.open_position_snapshots.is_empty());
5584 assert!(result.pending_order_snapshots.is_empty());
5585 }
5586
5587 #[test]
5588 fn exact_time_close_cannot_beat_later_symbol_stoploss() {
5589 let open_ts = ts(10, 0, 0);
5590 let execution_ts = ts(10, 0, 1);
5591 let events = vec![
5592 FeedEvent::new(
5593 tick("XAUUSD", 100.0, 100.0, open_ts),
5594 EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
5595 ),
5596 FeedEvent::new(
5597 tick("EURUSD", 1.1, 1.1, execution_ts),
5598 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
5599 ),
5600 FeedEvent::new(
5601 tick("XAUUSD", 98.0, 98.0, execution_ts),
5602 EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
5603 ),
5604 ];
5605 let signals = vec![
5606 RawSignal::Entry {
5607 ts: open_ts,
5608 symbol: "XAUUSD".into(),
5609 side: Side::Buy,
5610 order_type: OrderType::Market,
5611 price: None,
5612 risk_multiplier: 1.0,
5613 stoploss: Some(99.0),
5614 targets: Vec::new(),
5615 group: None,
5616 trade_id: Some("later-stop".into()),
5617 entry_class: None,
5618 },
5619 RawSignal::Close {
5620 ts: execution_ts,
5621 position: PositionRef::ByTradeId {
5622 trade_id: "later-stop".into(),
5623 },
5624 },
5625 ];
5626 let mut feed = VecFeed::from_feed_events(events);
5627 let result = BacktestRunner::new_future(
5628 BacktestConfig {
5629 close_on_finish: false,
5630 ..fixed_lot_config()
5631 },
5632 FutureQuoteConfig::default(),
5633 )
5634 .run_raw_signals_future(&mut feed, signals, None);
5635
5636 assert_eq!(result.close_events.len(), 1);
5637 assert_eq!(result.close_events[0].reason, CloseReason::Stoploss);
5638 assert_eq!(
5639 result.recorded_fills.last().unwrap().fill.purpose,
5640 FillPurpose::StopLoss
5641 );
5642 assert!(!result.action_dispositions.iter().any(|disposition| {
5643 disposition.action_id.starts_with("signal:00000001")
5644 && disposition.status == crate::ledger::ActionDispositionStatus::Applied
5645 }));
5646 }
5647
5648 #[test]
5649 fn exact_time_multisymbol_closes_preserve_signal_order() {
5650 let open_ts = ts(10, 0, 0);
5651 let close_ts = ts(10, 0, 1);
5652 let mut events = Vec::new();
5653 for (timestamp, row) in [(open_ts, 0), (close_ts, 1)] {
5654 events.push(FeedEvent::new(
5655 tick("EURUSD", 1.1, 1.1, timestamp),
5656 EventMetadata::new(SeriesRoles::PRIMARY, 0, row),
5657 ));
5658 events.push(FeedEvent::new(
5659 tick("XAUUSD", 100.0, 100.0, timestamp),
5660 EventMetadata::new(SeriesRoles::PRIMARY, 1, row),
5661 ));
5662 }
5663 let entry = |symbol: &str, trade_id: &str| RawSignal::Entry {
5664 ts: open_ts,
5665 symbol: symbol.into(),
5666 side: Side::Buy,
5667 order_type: OrderType::Market,
5668 price: None,
5669 risk_multiplier: 1.0,
5670 stoploss: None,
5671 targets: Vec::new(),
5672 group: None,
5673 trade_id: Some(trade_id.into()),
5674 entry_class: None,
5675 };
5676 let close = |trade_id: &str| RawSignal::Close {
5677 ts: close_ts,
5678 position: PositionRef::ByTradeId {
5679 trade_id: trade_id.into(),
5680 },
5681 };
5682 let signals = vec![
5683 entry("XAUUSD", "close-first"),
5684 entry("EURUSD", "close-second"),
5685 close("close-first"),
5686 close("close-second"),
5687 ];
5688 let mut feed = VecFeed::from_feed_events(events);
5689 let result = BacktestRunner::new_future(
5690 BacktestConfig {
5691 close_on_finish: false,
5692 ..fixed_lot_config()
5693 },
5694 FutureQuoteConfig::default(),
5695 )
5696 .run_raw_signals_future(&mut feed, signals, None);
5697
5698 assert_eq!(
5699 result
5700 .close_events
5701 .iter()
5702 .map(|event| event.symbol.as_str())
5703 .collect::<Vec<_>>(),
5704 vec!["XAUUSD", "EURUSD"]
5705 );
5706 assert!(
5707 result
5708 .close_events
5709 .iter()
5710 .all(|event| event.reason == CloseReason::Manual)
5711 );
5712 }
5713
5714 #[test]
5715 fn future_streaming_quiescence_waits_for_all_blockers() {
5716 let run = |events: Vec<MarketEvent>, signals: Vec<RawSignal>, config: BacktestConfig| {
5717 let polls = std::rc::Rc::new(std::cell::Cell::new(0));
5718 let primary_eod = events.last().map(MarketEvent::ts);
5719 let mut feed = CountingBatchFeed {
5720 batches: events.into_iter().map(primary_batch).collect(),
5721 polls: polls.clone(),
5722 };
5723 BacktestRunner::new_future(config, FutureQuoteConfig::default())
5724 .run_raw_signals_future_streaming_controlled(
5725 &mut feed,
5726 primary_eod,
5727 signals,
5728 None,
5729 || false,
5730 |_| {},
5731 )
5732 .unwrap();
5733 polls.get()
5734 };
5735 let eur_events = vec![
5736 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5737 tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
5738 tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
5739 ];
5740
5741 let immediately_quiescent = run(
5742 eur_events.clone(),
5743 vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
5744 BacktestConfig::default(),
5745 );
5746 assert_eq!(immediately_quiescent, 1);
5747
5748 let scheduled = run(
5749 eur_events.clone(),
5750 vec![RawSignal::CloseAll { ts: ts(10, 0, 2) }],
5751 BacktestConfig::default(),
5752 );
5753 assert_eq!(scheduled, 3, "scheduled signals must block termination");
5754
5755 let mut two_symbol_config = BacktestConfig {
5756 close_on_finish: false,
5757 ..fixed_lot_config()
5758 };
5759 two_symbol_config
5760 .symbol_specs
5761 .insert("GBPUSD".into(), test_symbol_spec("GBPUSD"));
5762 let queued = run(
5763 vec![
5764 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5765 tick("GBPUSD", 1.2500, 1.2502, ts(10, 0, 1)),
5766 tick("GBPUSD", 1.2501, 1.2503, ts(10, 0, 2)),
5767 ],
5768 vec![
5769 market_entry(ts(10, 0, 0), "GBPUSD", OrderType::Market),
5770 RawSignal::CloseAll { ts: ts(10, 0, 1) },
5771 ],
5772 two_symbol_config,
5773 );
5774 assert_eq!(
5775 queued, 2,
5776 "queued actions must wait for an eligible symbol quote"
5777 );
5778
5779 let open = run(
5780 eur_events.clone(),
5781 vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market)],
5782 BacktestConfig {
5783 close_on_finish: false,
5784 ..fixed_lot_config()
5785 },
5786 );
5787 assert_eq!(
5788 open, 4,
5789 "open positions must consume the stream through EOD"
5790 );
5791
5792 let pending = run(
5793 eur_events,
5794 vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Limit)],
5795 fixed_lot_config(),
5796 );
5797 assert_eq!(
5798 pending, 4,
5799 "pending orders must consume the stream through EOD"
5800 );
5801 }
5802
5803 #[test]
5804 fn future_mtm_output_policies_bound_curve_and_validate_before_feed_use() {
5805 assert_eq!(
5806 FutureQuoteConfig::default().mtm_output,
5807 MtmOutputPolicy::Bounded { max_points: 4_096 }
5808 );
5809 let events: Vec<_> = (0..12)
5810 .map(|second| tick("EURUSD", 100.0, 100.0, ts(10, 0, second)))
5811 .collect();
5812
5813 let pending = RawSignal::Entry {
5814 ts: ts(10, 0, 0),
5815 symbol: "EURUSD".into(),
5816 side: Side::Buy,
5817 order_type: OrderType::Limit,
5818 price: Some(90.0),
5819 risk_multiplier: 1.0,
5820 stoploss: None,
5821 targets: Vec::new(),
5822 group: None,
5823 trade_id: Some("mtm-policy-blocker".into()),
5824 entry_class: None,
5825 };
5826 let run = |policy| {
5827 let mut feed = VecFeed::new(events.clone());
5828 BacktestRunner::new_future(
5829 fixed_lot_config(),
5830 FutureQuoteConfig {
5831 mtm_output: policy,
5832 ..FutureQuoteConfig::default()
5833 },
5834 )
5835 .run_raw_signals_future(&mut feed, vec![pending.clone()], None)
5836 };
5837
5838 let none = run(MtmOutputPolicy::None);
5839 assert!(none.mtm_equity_curve.is_empty());
5840 assert_eq!(none.mtm_output_summary.observed_points, 13);
5841 assert_eq!(none.mtm_output_summary.omitted_points, 13);
5842
5843 let bounded = run(MtmOutputPolicy::Bounded { max_points: 8 });
5844 assert_eq!(bounded.mtm_equity_curve.len(), 8);
5845 assert_eq!(bounded.mtm_output_summary.observed_points, 13);
5846 assert_eq!(bounded.mtm_output_summary.retained_points, 8);
5847 assert_eq!(bounded.mtm_output_summary.omitted_points, 5);
5848
5849 let full = run(MtmOutputPolicy::Full);
5850 assert_eq!(full.mtm_equity_curve.len(), 13);
5851 assert_eq!(full.mtm_output_summary.observed_points, 13);
5852 assert_eq!(full.mtm_output_summary.omitted_points, 0);
5853 assert_eq!(
5854 full.mtm_equity_curve
5855 .iter()
5856 .filter(|point| {
5857 point.observation_kind.as_deref()
5858 == Some(EquityObservationKind::PostOutput.as_str())
5859 })
5860 .count(),
5861 0
5862 );
5863
5864 let mut invalid_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5865 let rejected = BacktestRunner::new_future(
5866 BacktestConfig::default(),
5867 FutureQuoteConfig {
5868 mtm_output: MtmOutputPolicy::Bounded { max_points: 7 },
5869 ..FutureQuoteConfig::default()
5870 },
5871 )
5872 .run_raw_signals_future(&mut invalid_feed, Vec::new(), None);
5873 assert_eq!(invalid_feed.remaining(), 1);
5874 assert!(rejected.action_dispositions.iter().any(|disposition| {
5875 disposition.action_id == "configuration"
5876 && disposition
5877 .reason
5878 .as_deref()
5879 .is_some_and(|reason| reason.contains("MTM max_points"))
5880 }));
5881 }
5882
5883 #[test]
5884 fn future_mtm_records_changed_post_output_observation_kind() {
5885 let mut feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5886 let signal = RawSignal::Entry {
5887 ts: ts(10, 0, 0),
5888 symbol: "EURUSD".into(),
5889 side: Side::Buy,
5890 order_type: OrderType::Market,
5891 price: None,
5892 risk_multiplier: 1.0,
5893 stoploss: None,
5894 targets: Vec::new(),
5895 group: None,
5896 trade_id: Some("mtm-kind".into()),
5897 entry_class: None,
5898 };
5899 let result = BacktestRunner::new_future(
5900 BacktestConfig {
5901 close_on_finish: false,
5902 ..fixed_lot_config()
5903 },
5904 FutureQuoteConfig {
5905 mtm_output: MtmOutputPolicy::Full,
5906 ..FutureQuoteConfig::default()
5907 },
5908 )
5909 .run_raw_signals_future(&mut feed, vec![signal], None);
5910
5911 let kinds: Vec<_> = result
5912 .mtm_equity_curve
5913 .iter()
5914 .filter_map(|point| point.observation_kind.as_deref())
5915 .collect();
5916 assert_eq!(
5917 kinds,
5918 vec![
5919 EquityObservationKind::PreSettlement.as_str(),
5920 EquityObservationKind::PostOutput.as_str(),
5921 EquityObservationKind::EndOfData.as_str(),
5922 ]
5923 );
5924 assert_eq!(
5925 result
5926 .execution_metadata
5927 .as_ref()
5928 .unwrap()
5929 .tags
5930 .get("termination_reason")
5931 .map(String::as_str),
5932 Some("end_of_data")
5933 );
5934 }
5935
5936 #[test]
5937 fn future_fallible_batch_feed_propagates_source_error() {
5938 let batch = TimestampBatch {
5939 ts: ts(10, 0, 0),
5940 events: vec![FeedEvent::new(
5941 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
5942 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
5943 )],
5944 };
5945 let mut feed = ScriptedBatchFeed {
5946 batches: VecDeque::from([Ok(Some(batch)), Err("feed failed")]),
5947 };
5948 let result =
5949 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
5950 .run_raw_signals_future_fallible(&mut feed, Vec::new(), None);
5951
5952 assert!(matches!(result, Err("feed failed")));
5953 }
5954
5955 struct BuyOnceStrategy {
5959 entered: bool,
5960 }
5961
5962 impl BuyOnceStrategy {
5963 fn new() -> Self {
5964 Self { entered: false }
5965 }
5966 }
5967
5968 impl Strategy for BuyOnceStrategy {
5969 fn on_event(&mut self, event: &MarketEvent) -> Vec<Action> {
5970 if self.entered {
5971 return vec![];
5972 }
5973 if let MarketEvent::Tick { symbol, ask, .. } = event {
5974 self.entered = true;
5975 vec![Action::Open {
5976 symbol: symbol.clone(),
5977 side: Side::Buy,
5978 order_type: OrderType::Market,
5979 price: Some(*ask),
5980 size: 1.0,
5981 stoploss: Some(*ask - 0.0050),
5982 targets: vec![TargetSpec {
5983 price: *ask + 0.0050,
5984 close_ratio: 1.0,
5985 }],
5986 rules: vec![],
5987 group: None,
5988 trade_id: None,
5989 }]
5990 } else {
5991 vec![]
5992 }
5993 }
5994
5995 fn on_finished(&mut self) -> Vec<Action> {
5996 vec![]
5999 }
6000 }
6001
6002 #[test]
6005 fn strategy_backtest_tp_hit() {
6006 let events = vec![
6007 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6008 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6009 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6010 tick("EURUSD", 1.0890, 1.0892, ts(10, 0, 3)),
6011 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 4)),
6013 ];
6014 let mut feed = VecFeed::new(events);
6015 let mut strategy = BuyOnceStrategy::new();
6016
6017 let config = BacktestConfig {
6018 initial_balance: 10_000.0,
6019 close_on_finish: true,
6020 ..Default::default()
6021 };
6022 let runner = BacktestRunner::new(config);
6023 let result = runner.run_strategy(&mut feed, &mut strategy);
6024
6025 assert_eq!(result.total_trades, 1);
6026 assert_eq!(result.winning_trades, 1);
6027 assert!(result.total_pnl > 0.0);
6028 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6029 }
6030
6031 #[test]
6032 fn strategy_backtest_sl_hit() {
6033 let events = vec![
6034 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6035 tick("EURUSD", 1.0830, 1.0832, ts(10, 0, 1)),
6036 tick("EURUSD", 1.0799, 1.0801, ts(10, 0, 2)),
6038 ];
6039 let mut feed = VecFeed::new(events);
6040 let mut strategy = BuyOnceStrategy::new();
6041
6042 let runner = BacktestRunner::with_defaults();
6043 let result = runner.run_strategy(&mut feed, &mut strategy);
6044
6045 assert_eq!(result.total_trades, 1);
6046 assert_eq!(result.losing_trades, 1);
6047 assert!(result.total_pnl < 0.0);
6048 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6049 }
6050
6051 #[test]
6052 fn strategy_close_on_finish() {
6053 let events = vec![
6055 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6056 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6057 tick("EURUSD", 1.0852, 1.0854, ts(10, 0, 2)),
6058 ];
6059 let mut feed = VecFeed::new(events);
6060 let mut strategy = BuyOnceStrategy::new();
6061
6062 let config = BacktestConfig {
6063 initial_balance: 10_000.0,
6064 close_on_finish: true,
6065 ..Default::default()
6066 };
6067 let runner = BacktestRunner::new(config);
6068 let result = runner.run_strategy(&mut feed, &mut strategy);
6069
6070 assert_eq!(result.total_trades, 1);
6071 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6072 }
6073
6074 #[test]
6075 fn strategy_no_close_on_finish() {
6076 let events = vec![
6077 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6078 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6079 ];
6080 let mut feed = VecFeed::new(events);
6081 let mut strategy = BuyOnceStrategy::new();
6082
6083 let config = BacktestConfig {
6084 initial_balance: 10_000.0,
6085 close_on_finish: false,
6086 ..Default::default()
6087 };
6088 let runner = BacktestRunner::new(config);
6089 let result = runner.run_strategy(&mut feed, &mut strategy);
6090
6091 assert_eq!(result.total_trades, 0);
6093 }
6094
6095 #[test]
6098 fn legacy_unprofiled_targets_default_to_equal_weights() {
6099 let events = vec![
6100 tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
6101 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
6102 tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
6103 ];
6104 let mut feed = VecFeed::new(events);
6105 let signals = vec![RawSignal::Entry {
6106 ts: ts(10, 0, 0),
6107 symbol: "EURUSD".into(),
6108 side: Side::Buy,
6109 order_type: OrderType::Market,
6110 price: Some(1.0000),
6111 risk_multiplier: 1.0,
6112 stoploss: None,
6113 targets: vec![1.1000, 1.2000],
6114 group: None,
6115 trade_id: Some("equal-targets".into()),
6116 entry_class: None,
6117 }];
6118
6119 let result = BacktestRunner::new(BacktestConfig {
6120 close_on_finish: false,
6121 ..fixed_lot_config()
6122 })
6123 .run_raw_signals(&mut feed, signals, None);
6124
6125 assert_eq!(result.trade_log.len(), 2);
6126 assert!(
6127 result
6128 .trade_log
6129 .iter()
6130 .all(|trade| (trade.size - 0.5).abs() < f64::EPSILON)
6131 );
6132 assert!(
6133 result
6134 .trade_log
6135 .iter()
6136 .all(|trade| trade.close_reason == CloseReason::Target)
6137 );
6138 }
6139
6140 #[test]
6141 fn legacy_atomic_target_modification_retains_profile_ratio() {
6142 let events = vec![
6143 tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
6144 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
6145 tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
6146 tick("EURUSD", 1.3000, 1.3000, ts(10, 0, 3)),
6147 ];
6148 let mut feed = VecFeed::new(events);
6149 let position = PositionRef::ByTradeId {
6150 trade_id: "modified-target".into(),
6151 };
6152 let signals = vec![
6153 RawSignal::Entry {
6154 ts: ts(10, 0, 0),
6155 symbol: "EURUSD".into(),
6156 side: Side::Buy,
6157 order_type: OrderType::Market,
6158 price: Some(1.0000),
6159 risk_multiplier: 1.0,
6160 stoploss: None,
6161 targets: vec![1.1000, 1.3000],
6162 group: None,
6163 trade_id: Some("modified-target".into()),
6164 entry_class: None,
6165 },
6166 RawSignal::ModifyTarget {
6167 ts: ts(10, 0, 0),
6168 position,
6169 old_price: 1.1000,
6170 new_price: 1.2000,
6171 },
6172 ];
6173 let profile = ManagementProfile {
6174 name: "non-default-ratios".into(),
6175 target_selection: None,
6176 use_targets: vec![1, 2],
6177 close_ratios: vec![0.25, 0.75],
6178 target_source: TargetSource::FromSignal,
6179 stoploss_mode: StoplossMode::FromSignal,
6180 rules: vec![],
6181 group_override: None,
6182 let_remainder_run: false,
6183 entry_geometry: EntryGeometryPolicy::Strict,
6184 };
6185
6186 let result = BacktestRunner::new(BacktestConfig {
6187 close_on_finish: false,
6188 ..fixed_lot_config()
6189 })
6190 .run_raw_signals(&mut feed, signals, Some(&profile));
6191
6192 assert_eq!(result.trade_log.len(), 2);
6193 assert!((result.trade_log[0].exit_price - 1.2000).abs() < f64::EPSILON);
6194 assert!((result.trade_log[0].size - 0.25).abs() < f64::EPSILON);
6195 assert!((result.trade_log[1].exit_price - 1.3000).abs() < f64::EPSILON);
6196 assert!((result.trade_log[1].size - 0.75).abs() < f64::EPSILON);
6197 }
6198
6199 #[test]
6200 fn run_raw_signals_entry_only() {
6201 let events = vec![
6202 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6203 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6204 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6205 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
6206 ];
6207 let mut feed = VecFeed::new(events);
6208
6209 let raw_signals = vec![RawSignal::Entry {
6210 ts: ts(10, 0, 0),
6211 symbol: "EURUSD".into(),
6212 side: Side::Buy,
6213 order_type: OrderType::Market,
6214 price: Some(1.0850),
6215 risk_multiplier: 1.0,
6216 stoploss: Some(1.0800),
6217 targets: vec![1.0900],
6218 group: None,
6219 trade_id: None,
6220 entry_class: None,
6221 }];
6222
6223 let runner = BacktestRunner::new(fixed_lot_config());
6224 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6225
6226 assert_eq!(result.total_trades, 1);
6227 assert_eq!(result.winning_trades, 1);
6228 }
6229
6230 #[test]
6231 fn run_raw_signals_open_then_close() {
6232 let events = vec![
6233 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6234 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6235 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6236 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6237 ];
6238 let mut feed = VecFeed::new(events);
6239
6240 let raw_signals = vec![
6241 RawSignal::Entry {
6242 ts: ts(10, 0, 0),
6243 symbol: "EURUSD".into(),
6244 side: Side::Buy,
6245 order_type: OrderType::Market,
6246 price: Some(1.0850),
6247 risk_multiplier: 1.0,
6248 stoploss: None,
6249 targets: vec![],
6250 group: None,
6251 trade_id: Some("t1".into()),
6252 entry_class: None,
6253 },
6254 RawSignal::Close {
6255 ts: ts(10, 0, 2),
6256 position: PositionRef::ByTradeId {
6257 trade_id: "t1".into(),
6258 },
6259 },
6260 ];
6261
6262 let config = BacktestConfig {
6263 initial_balance: 10_000.0,
6264 close_on_finish: false,
6265 ..fixed_lot_config()
6266 };
6267 let runner = BacktestRunner::new(config);
6268 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6269
6270 assert_eq!(result.total_trades, 1);
6271 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6272 }
6273
6274 #[test]
6275 fn run_raw_signals_open_then_modify_sl() {
6276 let events = vec![
6278 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6279 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
6280 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 2)),
6282 tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
6284 ];
6285 let mut feed = VecFeed::new(events);
6286
6287 let raw_signals = vec![
6288 RawSignal::Entry {
6289 ts: ts(10, 0, 0),
6290 symbol: "EURUSD".into(),
6291 side: Side::Buy,
6292 order_type: OrderType::Market,
6293 price: Some(1.0850),
6294 risk_multiplier: 1.0,
6295 stoploss: Some(1.0800),
6296 targets: vec![],
6297 group: None,
6298 trade_id: Some("t1".into()),
6299 entry_class: None,
6300 },
6301 RawSignal::ModifyStoploss {
6302 ts: ts(10, 0, 2),
6303 position: PositionRef::ByTradeId {
6304 trade_id: "t1".into(),
6305 },
6306 price: 1.0840,
6307 },
6308 ];
6309
6310 let config = BacktestConfig {
6311 initial_balance: 10_000.0,
6312 close_on_finish: true,
6313 ..fixed_lot_config()
6314 };
6315 let runner = BacktestRunner::new(config);
6316 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6317
6318 assert_eq!(result.total_trades, 1);
6319 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6320 }
6321
6322 #[test]
6323 fn run_raw_signals_open_then_partial_close() {
6324 let events = vec![
6325 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6326 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
6327 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6328 tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 3)),
6329 ];
6330 let mut feed = VecFeed::new(events);
6331
6332 let raw_signals = vec![
6333 RawSignal::Entry {
6334 ts: ts(10, 0, 0),
6335 symbol: "EURUSD".into(),
6336 side: Side::Buy,
6337 order_type: OrderType::Market,
6338 price: Some(1.0850),
6339 risk_multiplier: 1.0,
6340 stoploss: None,
6341 targets: vec![],
6342 group: None,
6343 trade_id: Some("t1".into()),
6344 entry_class: None,
6345 },
6346 RawSignal::ClosePartial {
6347 ts: ts(10, 0, 1),
6348 position: PositionRef::ByTradeId {
6349 trade_id: "t1".into(),
6350 },
6351 ratio: 0.5,
6352 },
6353 ];
6354
6355 let config = BacktestConfig {
6356 initial_balance: 10_000.0,
6357 close_on_finish: true,
6358 ..fixed_lot_config()
6359 };
6360 let runner = BacktestRunner::new(config);
6361 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6362
6363 assert!(result.total_trades >= 1);
6365 }
6366
6367 #[test]
6368 fn run_raw_signals_group_workflow() {
6369 let events = vec![
6370 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6371 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6372 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6373 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6374 tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 4)),
6375 ];
6376 let mut feed = VecFeed::new(events);
6377
6378 let raw_signals = vec![
6379 RawSignal::Entry {
6381 ts: ts(10, 0, 0),
6382 symbol: "EURUSD".into(),
6383 side: Side::Buy,
6384 order_type: OrderType::Market,
6385 price: Some(1.0850),
6386 risk_multiplier: 1.0,
6387 stoploss: None,
6388 targets: vec![],
6389 group: Some("grp1".into()),
6390 trade_id: Some("t1".into()),
6391 entry_class: None,
6392 },
6393 RawSignal::Entry {
6394 ts: ts(10, 0, 1),
6395 symbol: "EURUSD".into(),
6396 side: Side::Buy,
6397 order_type: OrderType::Market,
6398 price: Some(1.0857),
6399 risk_multiplier: 1.0,
6400 stoploss: None,
6401 targets: vec![],
6402 group: Some("grp1".into()),
6403 trade_id: Some("t2".into()),
6404 entry_class: None,
6405 },
6406 RawSignal::CloseAllInGroup {
6408 ts: ts(10, 0, 3),
6409 group_id: "grp1".into(),
6410 },
6411 ];
6412
6413 let config = BacktestConfig {
6414 initial_balance: 10_000.0,
6415 close_on_finish: false,
6416 ..fixed_lot_config()
6417 };
6418 let runner = BacktestRunner::new(config);
6419 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6420
6421 assert_eq!(result.total_trades, 2);
6422 }
6423
6424 #[test]
6425 fn run_raw_signals_close_all_of_symbol() {
6426 let events = vec![
6427 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6428 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6429 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6430 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6431 ];
6432 let mut feed = VecFeed::new(events);
6433
6434 let raw_signals = vec![
6435 RawSignal::Entry {
6436 ts: ts(10, 0, 0),
6437 symbol: "EURUSD".into(),
6438 side: Side::Buy,
6439 order_type: OrderType::Market,
6440 price: Some(1.0850),
6441 risk_multiplier: 1.0,
6442 stoploss: None,
6443 targets: vec![],
6444 group: None,
6445 trade_id: Some("t1".into()),
6446 entry_class: None,
6447 },
6448 RawSignal::Entry {
6449 ts: ts(10, 0, 0),
6450 symbol: "EURUSD".into(),
6451 side: Side::Buy,
6452 order_type: OrderType::Market,
6453 price: Some(1.0850),
6454 risk_multiplier: 0.5,
6455 stoploss: None,
6456 targets: vec![],
6457 group: None,
6458 trade_id: Some("t2".into()),
6459 entry_class: None,
6460 },
6461 RawSignal::CloseAllOf {
6462 ts: ts(10, 0, 2),
6463 symbol: "EURUSD".into(),
6464 },
6465 ];
6466
6467 let config = BacktestConfig {
6468 initial_balance: 10_000.0,
6469 close_on_finish: false,
6470 ..fixed_lot_config()
6471 };
6472 let runner = BacktestRunner::new(config);
6473 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6474
6475 assert_eq!(result.total_trades, 2);
6476 }
6477
6478 #[test]
6479 fn run_raw_signals_with_profile() {
6480 let events = vec![
6481 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6482 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6483 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
6484 ];
6485 let mut feed = VecFeed::new(events);
6486
6487 let profile = ManagementProfile {
6488 name: "test".into(),
6489 target_selection: None,
6490 use_targets: vec![1],
6491 close_ratios: vec![1.0],
6492 target_source: TargetSource::FromSignal,
6493 stoploss_mode: StoplossMode::FromSignal,
6494 rules: vec![],
6495 group_override: None,
6496 let_remainder_run: false,
6497 entry_geometry: EntryGeometryPolicy::Strict,
6498 };
6499
6500 let raw_signals = vec![RawSignal::Entry {
6501 ts: ts(10, 0, 0),
6502 symbol: "EURUSD".into(),
6503 side: Side::Buy,
6504 order_type: OrderType::Market,
6505 price: Some(1.0850),
6506 risk_multiplier: 1.0,
6507 stoploss: Some(1.0800),
6508 targets: vec![1.0900],
6509 group: None,
6510 trade_id: Some("t1".into()),
6511 entry_class: None,
6512 }];
6513
6514 let runner = BacktestRunner::new(fixed_lot_config());
6515 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6516
6517 assert_eq!(result.total_trades, 1);
6518 assert_eq!(result.winning_trades, 1);
6519 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6520 }
6521
6522 #[test]
6523 fn run_raw_signals_with_profile_preserves_trade_id() {
6524 let events = vec![
6527 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6528 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6529 ];
6530 let mut feed = VecFeed::new(events);
6531
6532 let profile = ManagementProfile {
6533 name: "test".into(),
6534 target_selection: None,
6535 use_targets: vec![1],
6536 close_ratios: vec![1.0],
6537 target_source: TargetSource::FromSignal,
6538 stoploss_mode: StoplossMode::FromSignal,
6539 rules: vec![],
6540 group_override: None,
6541 let_remainder_run: false,
6542 entry_geometry: EntryGeometryPolicy::Strict,
6543 };
6544
6545 let raw_signals = vec![
6546 RawSignal::Entry {
6547 ts: ts(10, 0, 0),
6548 symbol: "EURUSD".into(),
6549 side: Side::Buy,
6550 order_type: OrderType::Market,
6551 price: Some(1.0850),
6552 risk_multiplier: 1.0,
6553 stoploss: Some(1.0800),
6554 targets: vec![1.0900],
6555 group: None,
6556 trade_id: Some("msg-100".into()),
6557 entry_class: None,
6558 },
6559 RawSignal::Close {
6560 ts: ts(10, 0, 1),
6561 position: PositionRef::ByTradeId {
6562 trade_id: "msg-100".into(),
6563 },
6564 },
6565 ];
6566
6567 let runner = BacktestRunner::new(fixed_lot_config());
6568 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6569
6570 assert_eq!(result.total_trades, 1);
6571 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6572 }
6573
6574 #[test]
6575 fn run_raw_signals_no_profile() {
6576 let events = vec![
6578 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6579 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6580 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6581 ];
6582 let mut feed = VecFeed::new(events);
6583
6584 let raw_signals = vec![RawSignal::Entry {
6585 ts: ts(10, 0, 0),
6586 symbol: "EURUSD".into(),
6587 side: Side::Buy,
6588 order_type: OrderType::Market,
6589 price: Some(1.0850),
6590 risk_multiplier: 1.0,
6591 stoploss: None,
6592 targets: vec![],
6593 group: None,
6594 trade_id: None,
6595 entry_class: None,
6596 }];
6597
6598 let config = BacktestConfig {
6599 initial_balance: 10_000.0,
6600 close_on_finish: true,
6601 ..fixed_lot_config()
6602 };
6603 let runner = BacktestRunner::new(config);
6604 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6605
6606 assert_eq!(result.total_trades, 1);
6607 }
6608
6609 #[test]
6610 fn run_raw_signals_last_on_symbol_resolution() {
6611 let events = vec![
6613 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6614 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6615 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6616 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6617 ];
6618 let mut feed = VecFeed::new(events);
6619
6620 let raw_signals = vec![
6621 RawSignal::Entry {
6622 ts: ts(10, 0, 0),
6623 symbol: "EURUSD".into(),
6624 side: Side::Buy,
6625 order_type: OrderType::Market,
6626 price: Some(1.0850),
6627 risk_multiplier: 1.0,
6628 stoploss: None,
6629 targets: vec![],
6630 group: None,
6631 trade_id: Some("t1".into()),
6632 entry_class: None,
6633 },
6634 RawSignal::Entry {
6635 ts: ts(10, 0, 1),
6636 symbol: "EURUSD".into(),
6637 side: Side::Buy,
6638 order_type: OrderType::Market,
6639 price: Some(1.0857),
6640 risk_multiplier: 1.0,
6641 stoploss: None,
6642 targets: vec![],
6643 group: None,
6644 trade_id: Some("t2".into()),
6645 entry_class: None,
6646 },
6647 RawSignal::Close {
6649 ts: ts(10, 0, 2),
6650 position: PositionRef::ByTradeId {
6651 trade_id: "t2".into(),
6652 },
6653 },
6654 ];
6655
6656 let config = BacktestConfig {
6657 initial_balance: 10_000.0,
6658 close_on_finish: true,
6659 ..fixed_lot_config()
6660 };
6661 let runner = BacktestRunner::new(config);
6662 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6663
6664 assert_eq!(result.total_trades, 2);
6666 }
6667
6668 #[test]
6669 fn run_raw_signals_unresolved_ref_skipped() {
6670 let events = vec![
6672 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6673 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6674 ];
6675 let mut feed = VecFeed::new(events);
6676
6677 let raw_signals = vec![RawSignal::Close {
6678 ts: ts(10, 0, 0),
6679 position: PositionRef::ByTradeId {
6680 trade_id: "nonexistent".into(),
6681 },
6682 }];
6683
6684 let config = BacktestConfig {
6685 initial_balance: 10_000.0,
6686 close_on_finish: false,
6687 ..Default::default()
6688 };
6689 let runner = BacktestRunner::new(config);
6690 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6691
6692 assert_eq!(result.total_trades, 0);
6694 }
6695
6696 #[test]
6699 fn signal_replay_basic() {
6700 let events = vec![
6701 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6702 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6703 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6704 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
6706 ];
6707 let mut feed = VecFeed::new(events);
6708
6709 let raw_signals = vec![RawSignal::Entry {
6710 ts: ts(10, 0, 0),
6711 symbol: "EURUSD".into(),
6712 side: Side::Buy,
6713 order_type: OrderType::Market,
6714 price: Some(1.0850),
6715 risk_multiplier: 1.0,
6716 stoploss: Some(1.0800),
6717 targets: vec![1.0900],
6718 group: None,
6719 trade_id: None,
6720 entry_class: None,
6721 }];
6722
6723 let runner = BacktestRunner::new(fixed_lot_config());
6724 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6725
6726 assert_eq!(result.total_trades, 1);
6727 assert_eq!(result.winning_trades, 1);
6728 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6729 }
6730
6731 #[test]
6732 fn signal_replay_multiple_signals() {
6733 let events = vec![
6734 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6735 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6736 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
6738 tick("EURUSD", 1.0910, 1.0912, ts(10, 0, 3)),
6739 tick("EURUSD", 1.0920, 1.0922, ts(10, 0, 4)),
6740 ];
6741 let mut feed = VecFeed::new(events);
6742
6743 let raw_signals = vec![
6744 RawSignal::Entry {
6745 ts: ts(10, 0, 0),
6746 symbol: "EURUSD".into(),
6747 side: Side::Buy,
6748 order_type: OrderType::Market,
6749 price: Some(1.0850),
6750 risk_multiplier: 1.0,
6751 stoploss: Some(1.0800),
6752 targets: vec![1.0900],
6753 group: None,
6754 trade_id: Some("t1".into()),
6755 entry_class: None,
6756 },
6757 RawSignal::Entry {
6758 ts: ts(10, 0, 1),
6759 symbol: "EURUSD".into(),
6760 side: Side::Buy,
6761 order_type: OrderType::Market,
6762 price: Some(1.0857),
6763 risk_multiplier: 1.0,
6764 stoploss: None,
6765 targets: vec![],
6766 group: None,
6767 trade_id: Some("t2".into()),
6768 entry_class: None,
6769 },
6770 ];
6771
6772 let config = BacktestConfig {
6773 initial_balance: 10_000.0,
6774 close_on_finish: true,
6775 ..fixed_lot_config()
6776 };
6777 let runner = BacktestRunner::new(config);
6778 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6779
6780 assert!(result.total_trades >= 2);
6782 }
6783
6784 #[test]
6785 fn signal_replay_signal_before_data_filtered() {
6786 let events = vec![
6793 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6794 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6795 ];
6796 let mut feed = VecFeed::new(events);
6797
6798 let raw_signals = vec![RawSignal::Entry {
6799 ts: ts(9, 0, 0), symbol: "EURUSD".into(),
6801 side: Side::Buy,
6802 order_type: OrderType::Market,
6803 price: Some(1.0850),
6804 risk_multiplier: 1.0,
6805 stoploss: None,
6806 targets: vec![1.0900],
6807 group: None,
6808 trade_id: None,
6809 entry_class: None,
6810 }];
6811
6812 let runner = BacktestRunner::new(fixed_lot_config());
6813 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6814
6815 assert_eq!(result.total_trades, 1);
6816 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6817 }
6818
6819 #[test]
6820 fn empty_feed_empty_result() {
6821 let mut feed = VecFeed::new(vec![]);
6822 let mut strategy = BuyOnceStrategy::new();
6823
6824 let runner = BacktestRunner::with_defaults();
6825 let result = runner.run_strategy(&mut feed, &mut strategy);
6826
6827 assert_eq!(result.total_trades, 0);
6828 assert!((result.final_balance - 10_000.0).abs() < f64::EPSILON);
6829 }
6830
6831 #[test]
6832 fn report_display_does_not_panic() {
6833 let events = vec![
6834 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6835 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6836 ];
6837 let mut feed = VecFeed::new(events);
6838 let mut strategy = BuyOnceStrategy::new();
6839
6840 let runner = BacktestRunner::with_defaults();
6841 let result = runner.run_strategy(&mut feed, &mut strategy);
6842
6843 let _display = format!("{result}");
6844 }
6845
6846 #[test]
6847 fn run_raw_signals_with_profile_open_then_modify_sl_by_trade_id() {
6848 let events = vec![
6849 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6850 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6851 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6852 tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
6853 ];
6854 let mut feed = VecFeed::new(events);
6855
6856 let profile = ManagementProfile {
6857 name: "test".into(),
6858 target_selection: None,
6859 use_targets: vec![1],
6860 close_ratios: vec![1.0],
6861 target_source: TargetSource::FromSignal,
6862 stoploss_mode: StoplossMode::FromSignal,
6863 rules: vec![],
6864 group_override: None,
6865 let_remainder_run: false,
6866 entry_geometry: EntryGeometryPolicy::Strict,
6867 };
6868
6869 let raw_signals = vec![
6870 RawSignal::Entry {
6871 ts: ts(10, 0, 0),
6872 symbol: "EURUSD".into(),
6873 side: Side::Buy,
6874 order_type: OrderType::Market,
6875 price: Some(1.0850),
6876 risk_multiplier: 1.0,
6877 stoploss: Some(1.0800),
6878 targets: vec![1.0900],
6879 group: None,
6880 trade_id: Some("t1".into()),
6881 entry_class: None,
6882 },
6883 RawSignal::ModifyStoploss {
6884 ts: ts(10, 0, 2),
6885 position: PositionRef::ByTradeId {
6886 trade_id: "t1".into(),
6887 },
6888 price: 1.0840,
6889 },
6890 ];
6891
6892 let config = BacktestConfig {
6893 initial_balance: 10_000.0,
6894 close_on_finish: false,
6895 ..fixed_lot_config()
6896 };
6897 let runner = BacktestRunner::new(config);
6898 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6899
6900 assert_eq!(result.total_trades, 1);
6901 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6902 }
6903
6904 #[test]
6905 fn run_raw_signals_with_profile_open_then_close_partial_by_trade_id() {
6906 let events = vec![
6907 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6908 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6909 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6910 ];
6911 let mut feed = VecFeed::new(events);
6912
6913 let profile = ManagementProfile {
6914 name: "test".into(),
6915 target_selection: None,
6916 use_targets: vec![1],
6917 close_ratios: vec![1.0],
6918 target_source: TargetSource::FromSignal,
6919 stoploss_mode: StoplossMode::FromSignal,
6920 rules: vec![],
6921 group_override: None,
6922 let_remainder_run: false,
6923 entry_geometry: EntryGeometryPolicy::Strict,
6924 };
6925
6926 let raw_signals = vec![
6927 RawSignal::Entry {
6928 ts: ts(10, 0, 0),
6929 symbol: "EURUSD".into(),
6930 side: Side::Buy,
6931 order_type: OrderType::Market,
6932 price: Some(1.0850),
6933 risk_multiplier: 1.0,
6934 stoploss: Some(1.0800),
6935 targets: vec![1.0900],
6936 group: None,
6937 trade_id: Some("t1".into()),
6938 entry_class: None,
6939 },
6940 RawSignal::ClosePartial {
6941 ts: ts(10, 0, 1),
6942 position: PositionRef::ByTradeId {
6943 trade_id: "t1".into(),
6944 },
6945 ratio: 0.5,
6946 },
6947 ];
6948
6949 let config = BacktestConfig {
6950 initial_balance: 10_000.0,
6951 close_on_finish: false,
6952 ..fixed_lot_config()
6953 };
6954 let runner = BacktestRunner::new(config);
6955 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6956
6957 assert!(result.total_trades >= 1);
6959 }
6960
6961 #[test]
6962 fn run_raw_signals_multi_position_by_trade_id_with_profile() {
6963 let events = vec![
6967 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6968 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6969 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6970 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6971 ];
6972 let mut feed = VecFeed::new(events);
6973
6974 let profile = ManagementProfile {
6975 name: "test".into(),
6976 target_selection: None,
6977 use_targets: vec![1],
6978 close_ratios: vec![1.0],
6979 target_source: TargetSource::FromSignal,
6980 stoploss_mode: StoplossMode::FromSignal,
6981 rules: vec![],
6982 group_override: Some("alpha".into()),
6983 let_remainder_run: false,
6984 entry_geometry: EntryGeometryPolicy::Strict,
6985 };
6986
6987 let raw_signals = vec![
6988 RawSignal::Entry {
6989 ts: ts(10, 0, 0),
6990 symbol: "EURUSD".into(),
6991 side: Side::Buy,
6992 order_type: OrderType::Market,
6993 price: Some(1.0850),
6994 risk_multiplier: 1.0,
6995 stoploss: None,
6996 targets: vec![1.0910],
6997 group: None,
6998 trade_id: Some("t1".into()),
6999 entry_class: None,
7000 },
7001 RawSignal::Entry {
7002 ts: ts(10, 0, 1),
7003 symbol: "EURUSD".into(),
7004 side: Side::Buy,
7005 order_type: OrderType::Market,
7006 price: Some(1.0857),
7007 risk_multiplier: 1.0,
7008 stoploss: None,
7009 targets: vec![1.0910],
7010 group: None,
7011 trade_id: Some("t2".into()),
7012 entry_class: None,
7013 },
7014 RawSignal::Close {
7016 ts: ts(10, 0, 2),
7017 position: PositionRef::ByTradeId {
7018 trade_id: "t1".into(),
7019 },
7020 },
7021 ];
7022
7023 let config = BacktestConfig {
7024 initial_balance: 10_000.0,
7025 close_on_finish: true,
7026 ..fixed_lot_config()
7027 };
7028 let runner = BacktestRunner::new(config);
7029 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
7030
7031 assert_eq!(result.total_trades, 2);
7033 for trade in &result.trade_log {
7035 assert_eq!(trade.group.as_deref(), Some("alpha"));
7036 }
7037 }
7038
7039 #[test]
7040 fn merged_feed_manual_close_uses_correct_symbol_quote() {
7041 use crate::data_feed::MarketEvent;
7046 let events = vec![
7047 MarketEvent::Tick {
7048 symbol: "XAUUSD".into(),
7049 ts: ts(10, 0, 0),
7050 bid: 5000.0,
7051 ask: 5001.0,
7052 },
7053 MarketEvent::Tick {
7054 symbol: "GBPJPY".into(),
7055 ts: ts(10, 0, 1),
7056 bid: 210.0,
7057 ask: 211.0,
7058 },
7059 MarketEvent::Tick {
7060 symbol: "XAUUSD".into(),
7061 ts: ts(10, 0, 2),
7062 bid: 5050.0,
7063 ask: 5051.0,
7064 },
7065 MarketEvent::Tick {
7067 symbol: "GBPJPY".into(),
7068 ts: ts(10, 0, 3),
7069 bid: 212.0,
7070 ask: 213.0,
7071 },
7072 ];
7073 let mut feed = VecFeed::new(events);
7074
7075 let raw_signals = vec![
7076 RawSignal::Entry {
7077 ts: ts(10, 0, 0),
7078 symbol: "XAUUSD".into(),
7079 side: Side::Buy,
7080 order_type: OrderType::Market,
7081 price: Some(5000.0),
7082 risk_multiplier: 1.0,
7083 stoploss: None,
7084 targets: vec![],
7085 group: None,
7086 trade_id: Some("xau-1".into()),
7087 entry_class: None,
7088 },
7089 RawSignal::Close {
7091 ts: ts(10, 0, 3),
7092 position: PositionRef::ByTradeId {
7093 trade_id: "xau-1".into(),
7094 },
7095 },
7096 ];
7097
7098 let config = BacktestConfig {
7099 initial_balance: 10_000.0,
7100 close_on_finish: false,
7101 ..fixed_lot_config()
7102 };
7103 let runner = BacktestRunner::new(config);
7104 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
7105
7106 assert_eq!(result.total_trades, 1);
7107 let trade = &result.trade_log[0];
7108 assert_eq!(trade.symbol, "XAUUSD");
7109 assert!(
7111 trade.exit_price > 4000.0,
7112 "Exit price should be XAUUSD (~5050), got {}",
7113 trade.exit_price
7114 );
7115 }
7116
7117 fn long_tick_feed(count: usize) -> VecFeed {
7118 let start = ts(10, 0, 0);
7119 VecFeed::new(
7120 (0..count)
7121 .map(|index| {
7122 tick(
7123 "EURUSD",
7124 1.0848,
7125 1.0850,
7126 start + Duration::milliseconds(index as i64),
7127 )
7128 })
7129 .collect(),
7130 )
7131 }
7132
7133 #[test]
7134 fn legacy_replay_can_be_cancelled_during_event_processing() {
7135 let cancelled = std::cell::Cell::new(false);
7136 let mut feed = long_tick_feed(1_000);
7137 let outcome = BacktestRunner::with_defaults().run_raw_signals_controlled(
7138 &mut feed,
7139 Vec::new(),
7140 None,
7141 || cancelled.get(),
7142 |progress| {
7143 if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
7144 cancelled.set(true);
7145 }
7146 },
7147 );
7148
7149 assert_eq!(outcome.unwrap_err(), ReplayCancelled);
7150 assert!(
7151 feed.remaining() > 0,
7152 "cancellation must stop further replay"
7153 );
7154 }
7155
7156 #[test]
7157 fn future_quote_replay_can_be_cancelled_during_event_processing() {
7158 let cancelled = std::cell::Cell::new(false);
7159 let mut feed = long_tick_feed(1_000);
7160 let runner = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default());
7161 let pending = RawSignal::Entry {
7162 ts: ts(10, 0, 0),
7163 symbol: "EURUSD".into(),
7164 side: Side::Buy,
7165 order_type: OrderType::Limit,
7166 price: Some(1.0),
7167 risk_multiplier: 1.0,
7168 stoploss: None,
7169 targets: Vec::new(),
7170 group: None,
7171 trade_id: Some("cancellation-blocker".into()),
7172 entry_class: None,
7173 };
7174 let outcome = runner.run_raw_signals_controlled(
7175 &mut feed,
7176 vec![pending],
7177 None,
7178 || cancelled.get(),
7179 |progress| {
7180 if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
7181 cancelled.set(true);
7182 }
7183 },
7184 );
7185
7186 assert_eq!(outcome.unwrap_err(), ReplayCancelled);
7187 }
7188
7189 #[test]
7190 fn controlled_replay_progress_is_monotonic_and_reaches_event_total() {
7191 let mut feed = long_tick_feed(600);
7192 let mut updates = Vec::new();
7193 BacktestRunner::with_defaults()
7194 .run_raw_signals_controlled(
7195 &mut feed,
7196 Vec::new(),
7197 None,
7198 || false,
7199 |progress| updates.push(progress),
7200 )
7201 .unwrap();
7202
7203 assert!(updates.len() >= 3);
7204 assert!(updates.windows(2).all(|pair| {
7205 pair[0].processed_events <= pair[1].processed_events
7206 && pair[0].processed_signals <= pair[1].processed_signals
7207 && pair[0].total_events <= pair[1].total_events
7208 && pair[0].total_signals <= pair[1].total_signals
7209 }));
7210 assert_eq!(updates.last().unwrap().processed_events, 600);
7211 assert_eq!(updates.last().unwrap().total_events, 600);
7212 }
7213
7214 #[test]
7215 fn legacy_replay_skips_invalid_crossed_and_reversed_quotes_without_nonfinite_pnl() {
7216 let events = vec![
7217 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7218 tick("EURUSD", f64::NAN, 101.0, ts(10, 0, 1)),
7219 tick("EURUSD", 102.0, 101.0, ts(10, 0, 2)),
7220 tick("EURUSD", 90.0, 90.0, ts(9, 59, 59)),
7221 tick("EURUSD", 110.0, 110.0, ts(10, 0, 3)),
7222 ];
7223 let mut feed = VecFeed::new(events);
7224 let signals = vec![RawSignal::Entry {
7225 ts: ts(10, 0, 0),
7226 symbol: "EURUSD".into(),
7227 side: Side::Buy,
7228 order_type: OrderType::Market,
7229 price: Some(100.0),
7230 risk_multiplier: 1.0,
7231 stoploss: None,
7232 targets: vec![],
7233 group: None,
7234 trade_id: Some("safe-feed".into()),
7235 entry_class: None,
7236 }];
7237
7238 let result =
7239 BacktestRunner::new(fixed_lot_config()).run_raw_signals(&mut feed, signals, None);
7240 assert_eq!(result.trade_log.len(), 1);
7241 assert_eq!(result.trade_log[0].exit_price, 110.0);
7242 assert_eq!(result.trade_log[0].pnl, 10.0);
7243 assert!(result.total_pnl.is_finite());
7244 assert!(result.final_balance.is_finite());
7245 }
7246
7247 #[test]
7248 fn legacy_and_future_profile_replay_share_empty_ratio_target_resolution() {
7249 let profile = ManagementProfile {
7250 name: "equal-target".into(),
7251 target_selection: None,
7252 use_targets: vec![1],
7253 close_ratios: vec![],
7254 target_source: TargetSource::FromSignal,
7255 stoploss_mode: StoplossMode::FromSignal,
7256 rules: vec![],
7257 group_override: None,
7258 let_remainder_run: false,
7259 entry_geometry: EntryGeometryPolicy::Strict,
7260 };
7261 let signals = vec![RawSignal::Entry {
7262 ts: ts(10, 0, 0),
7263 symbol: "EURUSD".into(),
7264 side: Side::Buy,
7265 order_type: OrderType::Market,
7266 price: Some(100.0),
7267 risk_multiplier: 1.0,
7268 stoploss: None,
7269 targets: vec![101.0],
7270 group: None,
7271 trade_id: Some("profile-parity".into()),
7272 entry_class: None,
7273 }];
7274 let events = vec![
7275 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7276 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7277 ];
7278
7279 let mut legacy_feed = VecFeed::new(events.clone());
7280 let legacy = BacktestRunner::new(BacktestConfig {
7281 close_on_finish: false,
7282 ..fixed_lot_config()
7283 })
7284 .run_raw_signals(&mut legacy_feed, signals.clone(), Some(&profile));
7285 let mut future_feed = VecFeed::new(events);
7286 let future = BacktestRunner::new_future(
7287 BacktestConfig {
7288 close_on_finish: false,
7289 ..fixed_lot_config()
7290 },
7291 FutureQuoteConfig::default(),
7292 )
7293 .run_raw_signals_future(&mut future_feed, signals, Some(&profile));
7294
7295 assert_eq!(legacy.trade_log.len(), 1);
7296 assert_eq!(future.trade_log.len(), 1);
7297 assert_eq!(legacy.trade_log[0].close_reason, CloseReason::Target);
7298 assert_eq!(future.trade_log[0].close_reason, CloseReason::Target);
7299 assert_eq!(legacy.trade_log[0].size, future.trade_log[0].size);
7300 }
7301
7302 #[test]
7303 fn future_batch_sizes_from_shared_conversion_before_primary_and_uses_primary_eod() {
7304 let currency_plan = RunCurrencyPlan::new(
7305 "USD",
7306 ["EURUSD".to_owned()].into_iter().collect(),
7307 ["EURUSD".to_owned()].into_iter().collect(),
7308 [("EURUSD".to_owned(), "EUR".to_owned())]
7309 .into_iter()
7310 .collect(),
7311 [(
7312 "EUR".to_owned(),
7313 ConversionRoute::Direct {
7314 pair: FxPair {
7315 symbol: "EURUSD".to_owned(),
7316 base_currency: "EUR".to_owned(),
7317 quote_currency: "USD".to_owned(),
7318 },
7319 },
7320 )]
7321 .into_iter()
7322 .collect(),
7323 Vec::new(),
7324 )
7325 .unwrap();
7326 let mut config = fixed_lot_config();
7327 config.sizing = Some(SizingPolicy::FixedRiskAmount { amount: 12.0 });
7328 let future = FutureQuoteConfig {
7329 currency_plan: Some(currency_plan),
7330 conversion_stale_after_ms: 1_000,
7331 ..FutureQuoteConfig::default()
7332 };
7333 let events = vec![
7334 FeedEvent::new(
7335 tick("EURUSD", 1.1, 1.2, ts(10, 0, 0)),
7336 EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
7337 ),
7338 FeedEvent::new(
7339 tick("EURUSD", 2.0, 2.1, ts(10, 0, 1)),
7340 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
7341 ),
7342 ];
7343 let signals = vec![RawSignal::Entry {
7344 ts: ts(10, 0, 0),
7345 symbol: "EURUSD".into(),
7346 side: Side::Buy,
7347 order_type: OrderType::Market,
7348 price: Some(1.0),
7349 risk_multiplier: 1.0,
7350 stoploss: Some(1.19),
7351 targets: Vec::new(),
7352 group: None,
7353 trade_id: Some("shared-conversion".into()),
7354 entry_class: None,
7355 }];
7356
7357 let mut feed = VecFeed::from_feed_events(events);
7358 let result = BacktestRunner::new_future(config, future)
7359 .run_raw_signals_future(&mut feed, signals, None);
7360
7361 assert_eq!(result.recorded_fills.len(), 2);
7362 assert!((result.recorded_fills[0].fill.price - 1.2).abs() < 1.0e-12);
7363 assert!((result.recorded_fills[0].size - 10.0).abs() < 1.0e-12);
7364 assert_eq!(result.recorded_fills[1].execution_ts, Some(ts(10, 0, 0)));
7365 assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 0));
7366 assert!(
7367 result
7368 .mtm_equity_curve
7369 .iter()
7370 .all(|point| point.ts == ts(10, 0, 0))
7371 );
7372 }
7373
7374 #[test]
7375 fn conversion_only_batch_revalues_but_defers_execution_to_primary_quote() {
7376 let currency_plan = RunCurrencyPlan::new(
7377 "USD",
7378 ["EURUSD".to_owned()].into_iter().collect(),
7379 ["EURUSD".to_owned()].into_iter().collect(),
7380 [("EURUSD".to_owned(), "EUR".to_owned())]
7381 .into_iter()
7382 .collect(),
7383 [(
7384 "EUR".to_owned(),
7385 ConversionRoute::Direct {
7386 pair: FxPair {
7387 symbol: "EURUSD".to_owned(),
7388 base_currency: "EUR".to_owned(),
7389 quote_currency: "USD".to_owned(),
7390 },
7391 },
7392 )]
7393 .into_iter()
7394 .collect(),
7395 Vec::new(),
7396 )
7397 .unwrap();
7398 let config = BacktestConfig {
7399 close_on_finish: false,
7400 ..fixed_lot_config()
7401 };
7402 let future = FutureQuoteConfig {
7403 currency_plan: Some(currency_plan),
7404 conversion_stale_after_ms: 10_000,
7405 ..FutureQuoteConfig::default()
7406 };
7407 let events = vec![
7408 FeedEvent::new(
7409 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7410 EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
7411 ),
7412 FeedEvent::new(
7413 tick("EURUSD", 2.0, 2.0, ts(10, 0, 1)),
7414 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
7415 ),
7416 FeedEvent::new(
7417 tick("EURUSD", 110.0, 110.0, ts(10, 0, 2)),
7418 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
7419 ),
7420 ];
7421 let signals = vec![
7422 RawSignal::Entry {
7423 ts: ts(10, 0, 0),
7424 symbol: "EURUSD".into(),
7425 side: Side::Buy,
7426 order_type: OrderType::Market,
7427 price: None,
7428 risk_multiplier: 1.0,
7429 stoploss: None,
7430 targets: Vec::new(),
7431 group: None,
7432 trade_id: Some("conversion-only".into()),
7433 entry_class: None,
7434 },
7435 RawSignal::Close {
7436 ts: ts(10, 0, 1),
7437 position: PositionRef::ByTradeId {
7438 trade_id: "conversion-only".into(),
7439 },
7440 },
7441 ];
7442
7443 let mut feed = VecFeed::from_feed_events(events);
7444 let result = BacktestRunner::new_future(config, future)
7445 .run_raw_signals_future(&mut feed, signals, None);
7446
7447 assert_eq!(result.recorded_fills.len(), 2);
7448 assert_eq!(result.recorded_fills[0].quote_ts, ts(10, 0, 0));
7449 assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 2));
7450 assert!(
7451 result
7452 .mtm_equity_curve
7453 .iter()
7454 .any(|point| point.ts == ts(10, 0, 1))
7455 );
7456 assert_eq!(result.total_pnl, 20.0);
7457 assert_eq!(result.close_events[0].native_pnl, Some(10.0));
7458 assert_eq!(
7459 result.close_events[0]
7460 .pnl_conversion
7461 .as_ref()
7462 .unwrap()
7463 .operation_ts,
7464 ts(10, 0, 2)
7465 );
7466 }
7467
7468 #[test]
7469 fn exact_timestamp_close_updates_balance_before_later_risk_entry() {
7470 let mut config = fixed_lot_config();
7471 config.close_on_finish = false;
7472 config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
7473 let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
7474 spec.digits = 2;
7475 spec.pip_position = 2;
7476 spec.lot_base_units = 1;
7477 spec.lot_step_units = 1;
7478 let future = FutureQuoteConfig {
7479 currency_plan: Some(identity_currency_plan("EURUSD")),
7480 ..FutureQuoteConfig::default()
7481 };
7482 let signals = vec![
7483 RawSignal::Entry {
7484 ts: ts(10, 0, 0),
7485 symbol: "EURUSD".into(),
7486 side: Side::Buy,
7487 order_type: OrderType::Market,
7488 price: None,
7489 risk_multiplier: 1.0,
7490 stoploss: Some(99.0),
7491 targets: Vec::new(),
7492 group: None,
7493 trade_id: Some("first".into()),
7494 entry_class: None,
7495 },
7496 RawSignal::Close {
7497 ts: ts(10, 0, 1),
7498 position: PositionRef::ByTradeId {
7499 trade_id: "first".into(),
7500 },
7501 },
7502 RawSignal::Entry {
7503 ts: ts(10, 0, 1),
7504 symbol: "EURUSD".into(),
7505 side: Side::Buy,
7506 order_type: OrderType::Market,
7507 price: None,
7508 risk_multiplier: 1.0,
7509 stoploss: Some(100.0),
7510 targets: Vec::new(),
7511 group: None,
7512 trade_id: Some("second".into()),
7513 entry_class: None,
7514 },
7515 ];
7516 let mut feed = VecFeed::new(vec![
7517 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7518 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7519 ]);
7520
7521 let result = BacktestRunner::new_future(config, future)
7522 .run_raw_signals_future(&mut feed, signals, None);
7523
7524 assert!((result.total_pnl - 100.0).abs() < 1.0e-12);
7525 assert_eq!(result.open_position_snapshots.len(), 1);
7526 assert_eq!(
7527 result.open_position_snapshots[0].trade_id.as_deref(),
7528 Some("second")
7529 );
7530 assert!((result.open_position_snapshots[0].remaining_size - 101.0).abs() < 1.0e-12);
7531 }
7532
7533 #[test]
7534 fn pending_fill_keeps_placement_size_after_balance_changes() {
7535 let mut config = fixed_lot_config();
7536 config.close_on_finish = false;
7537 config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
7538 let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
7539 spec.digits = 2;
7540 spec.pip_position = 2;
7541 spec.lot_base_units = 1;
7542 spec.lot_step_units = 1;
7543 let future = FutureQuoteConfig {
7544 currency_plan: Some(identity_currency_plan("EURUSD")),
7545 market_entry_sizing_basis: MarketEntrySizingBasis::SignalEntryPrice,
7546 ..FutureQuoteConfig::default()
7547 };
7548 let signals = vec![
7549 RawSignal::Entry {
7550 ts: ts(10, 0, 0),
7551 symbol: "EURUSD".into(),
7552 side: Side::Buy,
7553 order_type: OrderType::Market,
7554 price: None,
7555 risk_multiplier: 1.0,
7556 stoploss: Some(99.0),
7557 targets: Vec::new(),
7558 group: None,
7559 trade_id: Some("market".into()),
7560 entry_class: None,
7561 },
7562 RawSignal::Entry {
7563 ts: ts(10, 0, 0),
7564 symbol: "EURUSD".into(),
7565 side: Side::Buy,
7566 order_type: OrderType::Limit,
7567 price: Some(99.0),
7568 risk_multiplier: 1.0,
7569 stoploss: Some(98.0),
7570 targets: Vec::new(),
7571 group: None,
7572 trade_id: Some("pending".into()),
7573 entry_class: None,
7574 },
7575 RawSignal::Close {
7576 ts: ts(10, 0, 1),
7577 position: PositionRef::ByTradeId {
7578 trade_id: "market".into(),
7579 },
7580 },
7581 ];
7582 let mut feed = VecFeed::new(vec![
7583 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7584 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7585 tick("EURUSD", 99.0, 99.0, ts(10, 0, 2)),
7586 ]);
7587
7588 let result = BacktestRunner::new_future(config, future)
7589 .run_raw_signals_future(&mut feed, signals, None);
7590
7591 assert_eq!(result.pending_order_snapshots.len(), 0);
7592 assert_eq!(result.open_position_snapshots.len(), 1);
7593 assert_eq!(
7594 result.open_position_snapshots[0].trade_id.as_deref(),
7595 Some("pending")
7596 );
7597 assert!((result.open_position_snapshots[0].remaining_size - 100.0).abs() < 1.0e-12);
7598 let metadata = result.execution_metadata.as_ref().unwrap();
7599 assert_eq!(metadata.market_entry_sizing.len(), 1);
7600 assert_eq!(
7601 metadata.market_entry_sizing[0].trade_id.as_deref(),
7602 Some("market")
7603 );
7604 assert_eq!(metadata.entry_profile_resolutions.len(), 2);
7605 assert!(metadata.entry_profile_resolutions.iter().any(|audit| {
7606 audit.trade_id.as_deref() == Some("pending")
7607 && audit.resolution_stage == EntryResolutionStage::PendingPlacement
7608 }));
7609 }
7610
7611 #[test]
7612 fn raw_entries_require_sizing_but_management_only_replay_does_not() {
7613 let entry = RawSignal::Entry {
7614 ts: ts(10, 0, 0),
7615 symbol: "EURUSD".into(),
7616 side: Side::Buy,
7617 order_type: OrderType::Market,
7618 price: None,
7619 risk_multiplier: 1.0,
7620 stoploss: None,
7621 targets: Vec::new(),
7622 group: None,
7623 trade_id: None,
7624 entry_class: None,
7625 };
7626 let mut entry_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
7627 let rejected =
7628 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
7629 .run_raw_signals_future(&mut entry_feed, vec![entry], None);
7630 assert!(rejected.action_dispositions.iter().any(|disposition| {
7631 disposition.action_id == "configuration"
7632 && disposition
7633 .reason
7634 .as_deref()
7635 .is_some_and(|reason| reason.contains("BacktestConfig.sizing"))
7636 }));
7637
7638 let mut management_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
7639 let management =
7640 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
7641 .run_raw_signals_future(
7642 &mut management_feed,
7643 vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
7644 None,
7645 );
7646 assert!(
7647 management
7648 .action_dispositions
7649 .iter()
7650 .all(|disposition| disposition.action_id != "configuration")
7651 );
7652 }
7653
7654 #[test]
7655 fn server_filter_signals_before_market_window() {
7656 let events = vec![
7662 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
7663 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
7664 ];
7665 let mut feed = VecFeed::new(events);
7666
7667 let raw_signals = vec![RawSignal::Entry {
7669 ts: NaiveDate::from_ymd_opt(2026, 1, 1)
7670 .unwrap()
7671 .and_hms_opt(0, 0, 0)
7672 .unwrap(),
7673 symbol: "EURUSD".into(),
7674 side: Side::Buy,
7675 order_type: OrderType::Market,
7676 price: Some(1.0850),
7677 risk_multiplier: 1.0,
7678 stoploss: None,
7679 targets: vec![1.0900],
7680 group: None,
7681 trade_id: None,
7682 entry_class: None,
7683 }];
7684
7685 let runner = BacktestRunner::new(fixed_lot_config());
7686 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
7687
7688 assert_eq!(result.total_trades, 1);
7690 }
7691
7692 #[test]
7693 fn future_replay_audits_profile_resolution_rejection() {
7694 let profile = ManagementProfile {
7695 name: "requires_stop".into(),
7696 target_selection: Some(crate::profile::TargetSelection::None),
7697 use_targets: vec![],
7698 close_ratios: vec![],
7699 target_source: TargetSource::FromSignal,
7700 stoploss_mode: StoplossMode::FromSignalDistance { multiplier: 1.5 },
7701 rules: vec![],
7702 group_override: None,
7703 let_remainder_run: true,
7704 entry_geometry: EntryGeometryPolicy::Strict,
7705 };
7706 let profiles = PreparedEntryProfiles::try_new(
7707 Some(profile),
7708 Vec::<(String, ManagementProfile)>::new(),
7709 )
7710 .unwrap();
7711 let signals = vec![RawSignal::Entry {
7712 ts: ts(10, 0, 0),
7713 symbol: "EURUSD".into(),
7714 side: Side::Buy,
7715 order_type: OrderType::Market,
7716 price: Some(1.1000),
7717 risk_multiplier: 1.0,
7718 stoploss: None,
7719 targets: vec![],
7720 group: None,
7721 trade_id: Some("missing-stop".into()),
7722 entry_class: None,
7723 }];
7724 let mut feed = VecFeed::new(vec![tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1))]);
7725 let result = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default())
7726 .with_entry_profiles(profiles)
7727 .run_raw_signals_future(&mut feed, signals, None);
7728 let audit = &result
7729 .execution_metadata
7730 .as_ref()
7731 .unwrap()
7732 .entry_profile_resolutions[0];
7733 assert_eq!(audit.outcome, ActionDispositionStatus::Rejected);
7734 assert_eq!(audit.rejection_stage.as_deref(), Some("profile_resolution"));
7735 assert!(audit.reason.as_deref().unwrap().contains("signal stoploss"));
7736 }
7737
7738 #[test]
7739 fn future_replay_routes_entry_profiles_and_audits_resolved_levels() {
7740 let default_profile = ManagementProfile {
7741 name: "default".into(),
7742 target_selection: Some(crate::profile::TargetSelection::None),
7743 use_targets: vec![],
7744 close_ratios: vec![],
7745 target_source: TargetSource::FromSignal,
7746 stoploss_mode: StoplossMode::FromSignal,
7747 rules: vec![],
7748 group_override: None,
7749 let_remainder_run: true,
7750 entry_geometry: EntryGeometryPolicy::Strict,
7751 };
7752 let expanded_profile = ManagementProfile {
7753 name: "expanded".into(),
7754 target_selection: None,
7755 use_targets: vec![],
7756 close_ratios: vec![1.0],
7757 target_source: TargetSource::StopDistanceMultiples {
7758 multiples: vec![1.0],
7759 },
7760 stoploss_mode: StoplossMode::FromSignalDistance { multiplier: 1.5 },
7761 rules: vec![],
7762 group_override: None,
7763 let_remainder_run: false,
7764 entry_geometry: EntryGeometryPolicy::Strict,
7765 };
7766 let profiles = PreparedEntryProfiles::try_new(
7767 Some(default_profile),
7768 [("expanded".to_owned(), expanded_profile)],
7769 )
7770 .unwrap();
7771 let signals = vec![
7772 RawSignal::Entry {
7773 ts: ts(10, 0, 0),
7774 symbol: "EURUSD".into(),
7775 side: Side::Buy,
7776 order_type: OrderType::Market,
7777 price: Some(1.1000),
7778 risk_multiplier: 1.0,
7779 stoploss: Some(1.0990),
7780 targets: vec![],
7781 group: Some("same-group".into()),
7782 trade_id: Some("default-entry".into()),
7783 entry_class: None,
7784 },
7785 RawSignal::Entry {
7786 ts: ts(10, 0, 1),
7787 symbol: "EURUSD".into(),
7788 side: Side::Buy,
7789 order_type: OrderType::Market,
7790 price: Some(1.1000),
7791 risk_multiplier: 1.0,
7792 stoploss: Some(1.0990),
7793 targets: vec![],
7794 group: Some("same-group".into()),
7795 trade_id: Some("expanded-entry".into()),
7796 entry_class: Some("expanded".into()),
7797 },
7798 ];
7799 let mut feed = VecFeed::new(vec![
7800 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
7801 tick("EURUSD", 1.1002, 1.1002, ts(10, 0, 2)),
7802 tick("EURUSD", 1.1003, 1.1003, ts(10, 0, 3)),
7803 ]);
7804 let result = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default())
7805 .with_entry_profiles(profiles)
7806 .run_raw_signals_future(&mut feed, signals, None);
7807
7808 let audits = &result
7809 .execution_metadata
7810 .as_ref()
7811 .unwrap()
7812 .entry_profile_resolutions;
7813 assert_eq!(audits.len(), 2);
7814 assert_eq!(
7815 audits[0].selection_source,
7816 EntryProfileSelectionSource::RunDefault
7817 );
7818 assert_eq!(audits[0].selected_profile_name.as_deref(), Some("default"));
7819 assert_eq!(
7820 audits[1].selection_source,
7821 EntryProfileSelectionSource::Mapped
7822 );
7823 assert_eq!(audits[1].entry_class.as_deref(), Some("expanded"));
7824 assert_eq!(audits[1].selected_profile_name.as_deref(), Some("expanded"));
7825 let levels = audits[1].level_resolution.as_ref().unwrap();
7826 let reference = audits[1].level_reference_price.unwrap();
7827 let stop = levels.resolved_stoploss.unwrap();
7828 let target = levels.resolved_targets[0];
7829 let risk_distance = reference - stop;
7830 assert!(risk_distance > 0.0);
7831 assert!((target - reference - risk_distance).abs() < 1.0e-9);
7832 }
7833}