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::Open {
4297 trade_id: Some(trade_id),
4298 ..
4299 } = &action
4300 && self.engine.manager.id_by_trade_id(trade_id).is_some()
4301 {
4302 let mut disposition = ActionDisposition::rejected(action_id, "duplicate_trade_id");
4303 disposition.action_kind = Some(action_kind);
4304 disposition.signal_ts = Some(signal_ts);
4305 disposition.effective_ts = Some(effective_ts);
4306 self.record_disposition(lifecycle, disposition);
4307 return false;
4308 }
4309 if let Action::ScaleIn { position_id, .. } = &action
4310 && future_executor.has_close(position_id)
4311 {
4312 let mut disposition =
4313 ActionDisposition::rejected(action_id, "scale_in_after_close_not_supported");
4314 disposition.action_kind = Some(action_kind);
4315 disposition.signal_ts = Some(signal_ts);
4316 disposition.effective_ts = Some(effective_ts);
4317 disposition.position_ids.push(position_id.clone());
4318 self.record_disposition(lifecycle, disposition);
4319 return false;
4320 }
4321
4322 let execution = match self.prepare_future_action(&mut action, execution, quote, pricer) {
4323 Ok(execution) => execution,
4324 Err(reason) => {
4325 let mut disposition = ActionDisposition::rejected(action_id, reason);
4326 disposition.action_kind = Some(action_kind);
4327 disposition.signal_ts = Some(signal_ts);
4328 disposition.effective_ts = Some(effective_ts);
4329 self.record_disposition(lifecycle, disposition);
4330 return false;
4331 }
4332 };
4333
4334 let engine_transaction = match execution {
4335 Some(execution) => self
4336 .engine
4337 .begin_priced_future_action(action, quote, execution),
4338 None => self.engine.begin_future_action(action, effective_ts),
4339 };
4340 let engine_transaction = match engine_transaction {
4341 Ok(transaction) => transaction,
4342 Err(error) => {
4343 let closed_state = matches!(
4346 error,
4347 FutureApplyError::Core(qs_core::CoreError::InvalidState { .. })
4348 ) && action_kind != "entry";
4349 let mut disposition = if closed_state {
4350 ActionDisposition::skipped(action_id, "position_closed")
4351 } else {
4352 ActionDisposition::rejected(action_id, error.to_string())
4353 };
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 committed_effects = engine_transaction.effects().to_vec();
4363 let mut affected = if FutureExecutor::requires_processing(engine_transaction.effects()) {
4364 match future_executor.process_future_effects_with_currency(
4365 engine_transaction.effects(),
4366 &self.engine,
4367 quote,
4368 Some(&action_id),
4369 Some(signal_ts),
4370 effective_ts,
4371 portfolio,
4372 Some(conversion_quotes),
4373 ) {
4374 Ok(affected) => affected,
4375 Err(error) => {
4376 engine_transaction.rollback(&mut self.engine);
4377 let mut disposition = ActionDisposition::failed(action_id, error.to_string());
4378 disposition.action_kind = Some(action_kind);
4379 disposition.signal_ts = Some(signal_ts);
4380 disposition.effective_ts = Some(effective_ts);
4381 self.record_disposition(lifecycle, disposition);
4382 return false;
4383 }
4384 }
4385 } else {
4386 Vec::new()
4387 };
4388 for future_effect in engine_transaction.effects() {
4389 match future_effect.effect() {
4390 Effect::OrderPlaced { id } | Effect::OrderCancelled { id } => {
4391 affected.push(id.clone());
4392 }
4393 _ => {}
4394 }
4395 }
4396 affected.sort();
4397 affected.dedup();
4398 let _ = engine_transaction.commit();
4399 self.record_committed_effects(committed_effects, Some(action_id.clone()));
4400
4401 let mut disposition = ActionDisposition::applied(action_id);
4402 disposition.action_kind = Some(action_kind);
4403 disposition.signal_ts = Some(signal_ts);
4404 disposition.effective_ts = Some(effective_ts);
4405 disposition.position_ids = affected;
4406 self.record_disposition(lifecycle, disposition);
4407 true
4408 }
4409
4410 fn prepare_future_action(
4411 &self,
4412 action: &mut Action,
4413 prepriced: Option<ExecutionFill>,
4414 quote: &PriceQuote,
4415 pricer: &ExecutionPricer,
4416 ) -> Result<Option<ExecutionFill>, String> {
4417 let mut execution = prepriced;
4418 match action {
4419 Action::Open {
4420 symbol,
4421 side,
4422 order_type,
4423 price,
4424 size,
4425 ..
4426 } => {
4427 if !valid_accounting_size(*size) {
4428 return Err(format!(
4429 "position size must be finite and greater than the accounting tolerance, got {size}"
4430 ));
4431 }
4432 if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
4433 return Err(format!(
4434 "supplied entry price must be finite and positive, got {price:?}"
4435 ));
4436 }
4437 if *order_type == OrderType::Market {
4438 let priced = match execution {
4439 Some(priced) => priced,
4440 None => pricer
4441 .market_entry(*side, quote, self.pip_size(symbol))
4442 .map_err(|error| error.to_string())?,
4443 };
4444 *price = Some(priced.price);
4445 execution = Some(priced);
4446 } else {
4447 if price.is_none() {
4448 return Err("pending entry requires a requested price".to_owned());
4449 }
4450 execution = None;
4451 }
4452 }
4453 Action::ScaleIn {
4454 position_id,
4455 price,
4456 size,
4457 ..
4458 } => {
4459 if !valid_accounting_size(*size) {
4460 return Err(format!(
4461 "scale-in size must be finite and greater than the accounting tolerance, got {size}"
4462 ));
4463 }
4464 if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
4465 return Err(format!(
4466 "supplied scale-in price must be finite and positive, got {price:?}"
4467 ));
4468 }
4469 let side = self
4470 .engine
4471 .get_position(position_id)
4472 .map(|position| position.data.side)
4473 .ok_or_else(|| format!("position not found: {position_id}"))?;
4474 let priced = match execution {
4475 Some(priced) => priced,
4476 None => pricer
4477 .market_entry(side, quote, self.pip_size("e.symbol))
4478 .map_err(|error| error.to_string())?,
4479 };
4480 *price = Some(priced.price);
4481 execution = Some(priced);
4482 }
4483 Action::ClosePosition { position_id } | Action::ClosePartial { position_id, .. } => {
4484 let position = self
4485 .engine
4486 .get_position(position_id)
4487 .ok_or_else(|| format!("position not found: {position_id}"))?;
4488 if position.data.symbol != quote.symbol {
4489 return Err(format!(
4490 "position symbol {} does not match quote symbol {}",
4491 position.data.symbol, quote.symbol
4492 ));
4493 }
4494 execution = Some(
4495 pricer
4496 .market_exit(
4497 position.data.side,
4498 quote,
4499 self.pip_size(&position.data.symbol),
4500 )
4501 .map_err(|error| error.to_string())?,
4502 );
4503 }
4504 _ => execution = None,
4505 }
4506 Ok(execution)
4507 }
4508
4509 fn prepare_triggering_pending(
4510 &self,
4511 quote: &PriceQuote,
4512 side: Option<Side>,
4513 pricer: &ExecutionPricer,
4514 ) -> (Vec<PreparedPendingFill>, Vec<(String, String)>) {
4515 let ids = self
4516 .engine
4517 .manager
4518 .pending_ids_by_symbol_sorted("e.symbol);
4519 let mut prepared = Vec::new();
4520 let mut failures = Vec::new();
4521 for id in ids {
4522 let Some(position) = self.engine.get_position(&id) else {
4523 continue;
4524 };
4525 if side.is_some_and(|side| position.data.side != side) {
4526 continue;
4527 }
4528 let Some(purpose) = position.pending_fill_purpose(quote, self.config.fill_model) else {
4529 continue;
4530 };
4531 let execution = match pricer.price(
4532 purpose,
4533 position.data.side,
4534 quote,
4535 position.data.pending_price,
4536 self.pip_size("e.symbol),
4537 ) {
4538 Ok(fill) => fill,
4539 Err(error) => {
4540 failures.push((id, error.to_string()));
4541 continue;
4542 }
4543 };
4544
4545 let size = position.data.size;
4546 if !valid_accounting_size(size) {
4547 failures.push((
4548 id,
4549 format!(
4550 "pending size must be finite and greater than the accounting tolerance, got {size}"
4551 ),
4552 ));
4553 continue;
4554 }
4555
4556 prepared.push(PreparedPendingFill {
4557 position_id: id,
4558 execution,
4559 size,
4560 });
4561 }
4562 (prepared, failures)
4563 }
4564
4565 fn pip_size(&self, symbol: &str) -> f64 {
4566 self.config
4567 .symbol_specs
4568 .get(symbol)
4569 .map(|spec| 10_f64.powi(-(spec.pip_position as i32)))
4570 .unwrap_or(0.0001)
4571 }
4572
4573 fn action_symbol(&self, action: &Action) -> Option<String> {
4574 match action {
4575 Action::Open { symbol, .. } => Some(symbol.clone()),
4576 Action::ClosePosition { position_id }
4577 | Action::ClosePartial { position_id, .. }
4578 | Action::ModifyStoploss { position_id, .. }
4579 | Action::MoveStoplossToEntry { position_id }
4580 | Action::AddTarget { position_id, .. }
4581 | Action::RemoveTarget { position_id, .. }
4582 | Action::ModifyTarget { position_id, .. }
4583 | Action::AddRule { position_id, .. }
4584 | Action::RemoveRule { position_id, .. }
4585 | Action::ScaleIn { position_id, .. }
4586 | Action::CancelPending { position_id } => self
4587 .engine
4588 .get_position(position_id)
4589 .map(|position| position.data.symbol.clone()),
4590 Action::CloseAllOf { symbol } | Action::ModifyAllStoploss { symbol, .. } => {
4591 Some(symbol.clone())
4592 }
4593 _ => None,
4594 }
4595 }
4596
4597 fn resolve_future_actions(&self, signal: &RawSignal) -> Vec<Action> {
4598 match signal {
4599 RawSignal::CloseAllOf { symbol, .. } => self
4600 .engine
4601 .manager
4602 .open_ids_by_symbol_sorted(symbol)
4603 .into_iter()
4604 .map(|position_id| Action::ClosePosition { position_id })
4605 .collect(),
4606 RawSignal::CloseAll { .. } => self
4607 .engine
4608 .manager
4609 .ids_by_status_sorted(PositionStatus::Open)
4610 .into_iter()
4611 .map(|position_id| Action::ClosePosition { position_id })
4612 .collect(),
4613 RawSignal::CloseAllInGroup { group_id, .. } => {
4614 let mut ids = self.engine.manager.open_ids_by_group(group_id);
4615 ids.sort();
4616 ids.into_iter()
4617 .map(|position_id| Action::ClosePosition { position_id })
4618 .collect()
4619 }
4620 RawSignal::CancelAllPending { .. } => self
4621 .engine
4622 .manager
4623 .ids_by_status_sorted(PositionStatus::Pending)
4624 .into_iter()
4625 .map(|position_id| Action::CancelPending { position_id })
4626 .collect(),
4627 _ => resolve_signal(signal, &self.engine),
4628 }
4629 }
4630}
4631
4632fn map_strategy_driver_error<FeedError, StrategyError>(
4633 error: StrategyDriverError<StrategyError>,
4634) -> StrategyReplayError<FeedError, StrategyError> {
4635 match error {
4636 StrategyDriverError::Series(error) => StrategyReplayError::Series(error),
4637 StrategyDriverError::SeriesView(error) => StrategyReplayError::SeriesView(error),
4638 StrategyDriverError::Analysis(error) => StrategyReplayError::Analysis(error),
4639 StrategyDriverError::Strategy(error) => StrategyReplayError::Strategy(error),
4640 StrategyDriverError::Runtime(error) => StrategyReplayError::Runtime(error),
4641 StrategyDriverError::WarmupSignals { timestamp } => {
4642 StrategyReplayError::WarmupSignals { timestamp }
4643 }
4644 StrategyDriverError::InvalidGeneratedSignal {
4645 signal_index,
4646 reason,
4647 } => StrategyReplayError::InvalidGeneratedSignal {
4648 signal_index,
4649 reason,
4650 },
4651 StrategyDriverError::TickExecutionRequired { symbol, timestamp } => {
4652 StrategyReplayError::TickExecutionRequired { symbol, timestamp }
4653 }
4654 }
4655}
4656
4657#[allow(clippy::too_many_arguments)]
4658fn observe_future_equity(
4659 portfolio: &mut PortfolioRecorder,
4660 future_executor: &FutureExecutor,
4661 ts: NaiveDateTime,
4662 conversion_quotes: &ConversionQuoteBook,
4663 kind: EquityObservationKind,
4664 collector: &mut MtmCurveCollector,
4665 last_candidate: &mut Option<EquityPoint>,
4666 suppress_unchanged_post_output: bool,
4667) {
4668 portfolio.set_realized_pnl(future_executor.realized_pnl());
4669 let mut point = portfolio.observe_with_currency(
4670 ts,
4671 future_executor.open_snapshots(),
4672 Some(conversion_quotes),
4673 );
4674 point.observation_kind = Some(kind.as_str().to_owned());
4675 if suppress_unchanged_post_output
4676 && last_candidate
4677 .as_ref()
4678 .is_some_and(|previous| same_equity_values(previous, &point))
4679 {
4680 return;
4681 }
4682 collector.observe(point.clone());
4683 *last_candidate = Some(point);
4684}
4685
4686fn same_equity_values(left: &EquityPoint, right: &EquityPoint) -> bool {
4687 let mut left = left.clone();
4688 let mut right = right.clone();
4689 left.observation_kind = None;
4690 left.observation_sequence = None;
4691 right.observation_kind = None;
4692 right.observation_sequence = None;
4693 left == right
4694}
4695
4696fn pending_fill_position_id(effect: &FutureEffect) -> Option<&str> {
4697 match effect.effect() {
4698 Effect::PositionOpened { id } => Some(id),
4699 _ => None,
4700 }
4701}
4702
4703fn valid_accounting_size(size: f64) -> bool {
4704 size.is_finite() && size > position_size_tolerance(size)
4705}
4706
4707fn explicit_instrument_spec<'a>(
4708 config: &'a BacktestConfig,
4709 symbol: &str,
4710) -> Option<&'a InstrumentSpec> {
4711 config
4712 .instrument_manifest
4713 .as_ref()?
4714 .instruments
4715 .get(symbol)
4716 .map(|artifact| &artifact.spec)
4717}
4718
4719fn decimal_to_f64(value: Decimal, field: &str) -> Result<f64, String> {
4720 let value = value
4721 .to_string()
4722 .parse::<f64>()
4723 .map_err(|error| format!("invalid {field}: {error}"))?;
4724 if value.is_finite() {
4725 Ok(value)
4726 } else {
4727 Err(format!("{field} must be finite"))
4728 }
4729}
4730
4731fn instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
4732 decimal_to_f64(
4733 spec.economics.contract_multiplier.get(),
4734 "instrument contract multiplier",
4735 )
4736 .and_then(|value| {
4737 if value > 0.0 {
4738 Ok(value)
4739 } else {
4740 Err("instrument contract multiplier must be positive".into())
4741 }
4742 })
4743}
4744
4745fn supported_instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
4746 if spec.status != ListingStatus::Trading {
4747 return Err(format!(
4748 "instrument {} is not in trading status",
4749 spec.instrument
4750 ));
4751 }
4752 if spec.economics.quantity_unit != QuantityUnit::StandardLot {
4753 return Err(format!(
4754 "unsupported quantity unit for instrument {}: {:?}",
4755 spec.instrument, spec.economics.quantity_unit
4756 ));
4757 }
4758 let model = spec.economics.pnl_model.as_str();
4759 if model != EconomicsModelId::FX_QUOTE_LINEAR_V1
4760 && model != EconomicsModelId::CFD_QUOTE_LINEAR_V1
4761 {
4762 return Err(format!(
4763 "unsupported P&L model for instrument {}: {model}",
4764 spec.instrument
4765 ));
4766 }
4767 instrument_multiplier(spec)
4768}
4769
4770fn validate_instrument_manifest(config: &BacktestConfig) -> Result<(), String> {
4771 let Some(manifest) = &config.instrument_manifest else {
4772 return Ok(());
4773 };
4774 for (symbol, artifact) in &manifest.instruments {
4775 if symbol.is_empty() {
4776 return Err("instrument manifest symbol must not be empty".into());
4777 }
4778 artifact
4779 .spec
4780 .validate()
4781 .map_err(|error| format!("invalid instrument spec for {symbol}: {error}"))?;
4782 if artifact.resolved.instrument != artifact.spec.instrument {
4783 return Err(format!(
4784 "resolved instrument and spec identity differ for {symbol}"
4785 ));
4786 }
4787 if artifact.resolved.spec_revision != artifact.spec.revision {
4788 return Err(format!(
4789 "resolved specification revision does not match the embedded spec for {symbol}"
4790 ));
4791 }
4792 supported_instrument_multiplier(&artifact.spec)?;
4793 }
4794 for binding in &manifest.stored_series {
4795 let known = manifest
4796 .instruments
4797 .values()
4798 .any(|artifact| artifact.resolved == binding.instrument);
4799 if !known {
4800 return Err(format!(
4801 "stored series {}:{} references an instrument outside the manifest",
4802 binding.source_partition, binding.source_symbol
4803 ));
4804 }
4805 let artifact = manifest
4806 .instruments
4807 .values()
4808 .find(|artifact| artifact.resolved == binding.instrument)
4809 .expect("known binding reference has an instrument artifact");
4810 if binding.effective != artifact.spec.effective {
4811 return Err(format!(
4812 "stored series {}:{} effective interval differs from its instrument spec",
4813 binding.source_partition, binding.source_symbol
4814 ));
4815 }
4816 }
4817 Ok(())
4818}
4819
4820fn effective_contract_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
4821 let mut contract_sizes = config.contract_sizes.clone();
4822 if let Some(manifest) = &config.instrument_manifest {
4823 for (symbol, artifact) in &manifest.instruments {
4824 if let Ok(multiplier) = instrument_multiplier(&artifact.spec) {
4825 contract_sizes.insert(symbol.clone(), multiplier);
4826 }
4827 }
4828 }
4829 contract_sizes
4830}
4831
4832fn validate_replay_costs(
4834 config: &BacktestConfig,
4835 future: Option<&FutureQuoteConfig>,
4836) -> Result<(), String> {
4837 let account_currency = future
4838 .and_then(|future| future.currency_plan.as_ref())
4839 .map(|plan| plan.account_currency().to_owned());
4840 for (symbol, costs) in &config.costs {
4841 if symbol.is_empty() {
4842 return Err("cost symbol must not be empty".into());
4843 }
4844 costs
4845 .validate()
4846 .map_err(|error| format!("costs for {symbol} are invalid: {error}"))?;
4847 if let Some(account_currency) = account_currency.as_deref() {
4848 costs
4849 .validate_against_account_currency(account_currency)
4850 .map_err(|error| format!("costs for {symbol} are invalid: {error}"))?;
4851 }
4852 if costs.requires_point_size() && !config.symbol_specs.contains_key(symbol) {
4853 return Err(format!(
4854 "point-denominated swap for {symbol} requires a symbol specification for its digit count"
4855 ));
4856 }
4857 }
4858 Ok(())
4859}
4860
4861fn effective_point_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
4863 config
4864 .symbol_specs
4865 .iter()
4866 .map(|(symbol, spec)| (symbol.clone(), 10f64.powi(-i32::from(spec.digits))))
4867 .collect()
4868}
4869
4870fn is_monetary_sizing(policy: &SizingPolicy) -> bool {
4871 matches!(
4872 policy,
4873 SizingPolicy::FixedRiskAmount { .. } | SizingPolicy::BalanceRiskPercent { .. }
4874 )
4875}
4876
4877fn accept_legacy_quote(
4878 quote: &PriceQuote,
4879 last_quote_ts: &mut BTreeMap<String, NaiveDateTime>,
4880) -> bool {
4881 if ExecutionPricer::validate_quote(quote).is_err()
4882 || last_quote_ts
4883 .get("e.symbol)
4884 .is_some_and(|last| *last > quote.ts)
4885 {
4886 return false;
4887 }
4888 last_quote_ts.insert(quote.symbol.clone(), quote.ts);
4889 true
4890}
4891
4892pub const MAX_RUN_TAGS: usize = 32;
4894pub const MAX_RUN_TAG_BYTES: usize = 64;
4896
4897fn validate_run_tags(tags: &BTreeMap<String, String>) -> Result<(), String> {
4898 if tags.len() > MAX_RUN_TAGS {
4899 return Err(format!(
4900 "run tags must not exceed {MAX_RUN_TAGS} entries, got {}",
4901 tags.len()
4902 ));
4903 }
4904 for (key, value) in tags {
4905 if key.is_empty() || key.len() > MAX_RUN_TAG_BYTES {
4906 return Err(format!(
4907 "run tag key must be 1 to {MAX_RUN_TAG_BYTES} bytes, got '{key}'"
4908 ));
4909 }
4910 if !key
4911 .chars()
4912 .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
4913 {
4914 return Err(format!(
4915 "run tag key must be ASCII alphanumeric, underscore, or hyphen, got '{key}'"
4916 ));
4917 }
4918 if value.len() > MAX_RUN_TAG_BYTES {
4919 return Err(format!(
4920 "run tag value for '{key}' must not exceed {MAX_RUN_TAG_BYTES} bytes"
4921 ));
4922 }
4923 if value.chars().any(char::is_control) {
4924 return Err(format!(
4925 "run tag value for '{key}' must not contain control characters"
4926 ));
4927 }
4928 }
4929 Ok(())
4930}
4931
4932fn validate_replay_config(
4933 config: &BacktestConfig,
4934 future: Option<&FutureQuoteConfig>,
4935 raw_signals: &[RawSignal],
4936) -> Result<(), String> {
4937 if !config.initial_balance.is_finite() || config.initial_balance <= 0.0 {
4938 return Err(format!(
4939 "initial balance must be finite and positive, got {}",
4940 config.initial_balance
4941 ));
4942 }
4943 for (symbol, contract_size) in &config.contract_sizes {
4944 if symbol.is_empty() {
4945 return Err("contract-size symbol must not be empty".into());
4946 }
4947 if !contract_size.is_finite() || *contract_size <= 0.0 {
4948 return Err(format!(
4949 "contract size for {symbol} must be finite and positive, got {contract_size}"
4950 ));
4951 }
4952 }
4953
4954 validate_run_tags(&config.run_tags)?;
4955 validate_replay_costs(config, future)?;
4956 validate_instrument_manifest(config)?;
4957 for (symbol, spec) in &config.symbol_specs {
4958 validate_symbol_spec(symbol, spec)?;
4959 if explicit_instrument_spec(config, symbol).is_none() {
4960 resolve_legacy_economics(spec).map_err(|error| error.to_string())?;
4961 }
4962 }
4963
4964 let entry_symbols: Vec<&str> = raw_signals
4965 .iter()
4966 .filter_map(|signal| match signal {
4967 RawSignal::Entry { symbol, .. } => Some(symbol.as_str()),
4968 _ => None,
4969 })
4970 .collect();
4971 if !entry_symbols.is_empty() && config.sizing.is_none() {
4972 return Err("raw entry requires BacktestConfig.sizing".to_owned());
4973 }
4974 if let Some(policy) = &config.sizing {
4975 validate_sizing_policy(policy)?;
4976 for symbol in &entry_symbols {
4977 if !config.symbol_specs.contains_key(*symbol)
4978 && explicit_instrument_spec(config, symbol).is_none()
4979 {
4980 return Err(format!("missing instrument or symbol spec for {symbol}"));
4981 }
4982 }
4983 if !entry_symbols.is_empty() && is_monetary_sizing(policy) {
4984 let future = future.ok_or_else(|| {
4985 "monetary sizing requires FutureQuote execution and a currency plan".to_owned()
4986 })?;
4987 let plan = future
4988 .currency_plan
4989 .as_ref()
4990 .ok_or_else(|| "monetary sizing requires a FutureQuote currency plan".to_owned())?;
4991 for symbol in &entry_symbols {
4992 if plan.route_for_primary_symbol(symbol).is_none() {
4993 return Err(format!(
4994 "currency plan has no frozen route for primary symbol {symbol}"
4995 ));
4996 }
4997 }
4998 }
4999 }
5000
5001 if let Some(future) = future {
5002 if future.signal_latency_ms < 0 {
5003 return Err(format!(
5004 "signal latency must be non-negative, got {}",
5005 future.signal_latency_ms
5006 ));
5007 }
5008 let latency = Duration::milliseconds(future.signal_latency_ms);
5009 for signal in raw_signals {
5010 if signal.ts().checked_add_signed(latency).is_none() {
5011 return Err(format!(
5012 "signal latency overflows datetime for signal at {}",
5013 signal.ts()
5014 ));
5015 }
5016 }
5017 if !future.slippage_pips.is_finite() {
5018 return Err(format!(
5019 "slippage pips must be finite, got {}",
5020 future.slippage_pips
5021 ));
5022 }
5023 if future.stale_quote_after_ms.is_some_and(|value| value < 0) {
5024 return Err("stale quote threshold must be non-negative".into());
5025 }
5026 if !future.pnl_epsilon.is_finite() || future.pnl_epsilon < 0.0 {
5027 return Err(format!(
5028 "P&L epsilon must be finite and non-negative, got {}",
5029 future.pnl_epsilon
5030 ));
5031 }
5032 if future.conversion_stale_after_ms < 0 {
5033 return Err("conversion quote threshold must be non-negative".to_owned());
5034 }
5035 future
5036 .mtm_output
5037 .validate()
5038 .map_err(|error| error.to_string())?;
5039 }
5040 Ok(())
5041}
5042
5043fn validate_sizing_policy(policy: &SizingPolicy) -> Result<(), String> {
5044 let (name, value) = match policy {
5045 SizingPolicy::FixedLot { lots } => ("fixed lots", *lots),
5046 SizingPolicy::FixedRiskAmount { amount } => ("fixed risk amount", *amount),
5047 SizingPolicy::BalanceRiskPercent { percent } => ("balance risk percent", *percent),
5048 };
5049 if value.is_finite() && value > 0.0 {
5050 Ok(())
5051 } else {
5052 Err(format!("{name} must be finite and positive, got {value}"))
5053 }
5054}
5055
5056fn validate_symbol_spec(symbol: &str, spec: &qs_symbols::SymbolSpec) -> Result<(), String> {
5057 if symbol.is_empty() || spec.canonical.is_empty() {
5058 return Err("symbol spec names must not be empty".into());
5059 }
5060 if spec.digits > 18 || spec.pip_position > spec.digits {
5061 return Err(format!(
5062 "invalid price precision for {symbol}: digits={}, pip_position={}",
5063 spec.digits, spec.pip_position
5064 ));
5065 }
5066 if spec.lot_base_units <= 0
5067 || spec.lot_step_units <= 0
5068 || spec.lot_min_steps <= 0
5069 || spec.lot_max_steps < 0
5070 || (spec.lot_max_steps > 0 && spec.lot_max_steps < spec.lot_min_steps)
5071 {
5072 return Err(format!("invalid lot metadata for {symbol}"));
5073 }
5074 let lot_step = spec.lot_step();
5075 let min_lot = spec.lot_min();
5076 let max_lot = spec.lot_max();
5077 if !lot_step.is_finite()
5078 || lot_step <= 0.0
5079 || !min_lot.is_finite()
5080 || min_lot <= 0.0
5081 || !max_lot.is_finite()
5082 {
5083 return Err(format!("invalid derived lot metadata for {symbol}"));
5084 }
5085 Ok(())
5086}
5087
5088fn rejected_legacy_result(config: &BacktestConfig) -> BacktestResult {
5089 BacktestResult::from_trade_log(
5090 if config.initial_balance.is_finite() {
5091 config.initial_balance
5092 } else {
5093 0.0
5094 },
5095 Vec::new(),
5096 )
5097}
5098
5099fn rejected_future_result(
5100 config: &BacktestConfig,
5101 future: &FutureQuoteConfig,
5102 evaluation_options: EvaluationOptions,
5103 error: String,
5104) -> BacktestResult {
5105 let execution_model = ExecutionModel::new(
5106 qs_core::types::ExecutionConvention::FutureQuoteV1,
5107 config.fill_model,
5108 if future.slippage_pips == 0.0 {
5109 SlippageModel::None
5110 } else {
5111 SlippageModel::FixedPips {
5112 pips: future.slippage_pips,
5113 }
5114 },
5115 );
5116 let mut lifecycle = LifecycleLedger::new();
5117 let _ = lifecycle.record(ActionDisposition::rejected(
5118 "configuration",
5119 format!("invalid_configuration: {error}"),
5120 ));
5121 let mut tags = BTreeMap::new();
5122 tags.insert("configuration_error".into(), error);
5123 insert_economic_support_metadata(&mut tags, config);
5124 let artifacts = FutureBacktestArtifacts {
5125 execution: ExecutionMetadata {
5126 execution_model,
5127 initial_balance: if config.initial_balance.is_finite() {
5128 config.initial_balance
5129 } else {
5130 0.0
5131 },
5132 account_currency: future
5133 .currency_plan
5134 .as_ref()
5135 .map(|plan| plan.account_currency().to_owned()),
5136 currency_plan: future.currency_plan.clone(),
5137 contract_sizes: effective_contract_sizes(config)
5138 .into_iter()
5139 .filter(|(symbol, size)| !symbol.is_empty() && size.is_finite() && *size > 0.0)
5140 .collect(),
5141 instrument_manifest: config.instrument_manifest.clone(),
5142 instrument_sizing: Vec::new(),
5143 market_entry_sizing_basis: future.market_entry_sizing_basis,
5144 market_entry_sizing: Vec::new(),
5145 stale_quote_after_millis: future.stale_quote_after_ms,
5146 pnl_epsilon: if future.pnl_epsilon.is_finite() && future.pnl_epsilon >= 0.0 {
5147 future.pnl_epsilon
5148 } else {
5149 crate::artifacts::DEFAULT_PNL_EPSILON
5150 },
5151 tags,
5152 ..ExecutionMetadata::default()
5153 },
5154 lifecycle,
5155 mtm_output_summary: MtmOutputSummary {
5156 policy: future.mtm_output,
5157 ..MtmOutputSummary::default()
5158 },
5159 ..FutureBacktestArtifacts::default()
5160 };
5161 BacktestResult::from_future_artifacts_with_options(artifacts, evaluation_options)
5162}
5163
5164fn insert_economic_support_metadata(tags: &mut BTreeMap<String, String>, config: &BacktestConfig) {
5165 let mut compatibility_specs = config.symbol_specs.iter().peekable();
5166 if compatibility_specs.peek().is_none() {
5167 return;
5168 }
5169 tags.insert(
5170 "economics.guard".into(),
5171 LEGACY_ECONOMIC_GUARD_ID.to_owned(),
5172 );
5173 for (symbol, spec) in compatibility_specs {
5174 let prefix = format!("economics.symbol.{symbol}");
5175 tags.insert(format!("{prefix}.category"), spec.category.clone());
5176 match resolve_legacy_economics(spec) {
5177 Ok(economics) => {
5178 tags.insert(format!("{prefix}.status"), "supported".into());
5179 tags.insert(format!("{prefix}.model"), economics.model.as_str().into());
5180 tags.insert(
5181 format!("{prefix}.contract_multiplier"),
5182 economics.contract_multiplier.to_string(),
5183 );
5184 }
5185 Err(error) => {
5186 tags.insert(format!("{prefix}.status"), "unsupported".into());
5187 tags.insert(format!("{prefix}.reason"), error.to_string());
5188 }
5189 }
5190 }
5191}
5192
5193fn queued_exposure_symbols(
5194 queued: &VecDeque<QueuedAction>,
5195 quotes: &BTreeMap<String, PriceQuote>,
5196 batch_ts: NaiveDateTime,
5197) -> BTreeSet<String> {
5198 queued
5199 .iter()
5200 .filter(|action| {
5201 action.effective_ts <= batch_ts
5202 && quotes.contains_key(&action.symbol)
5203 && is_exposure_increasing(&action.action)
5204 })
5205 .map(|action| action.symbol.clone())
5206 .collect()
5207}
5208
5209fn is_exposure_increasing(action: &Action) -> bool {
5210 matches!(
5211 action,
5212 Action::Open {
5213 order_type: OrderType::Market,
5214 ..
5215 } | Action::ScaleIn { .. }
5216 )
5217}
5218
5219fn is_fill_bearing(action: &Action) -> bool {
5220 matches!(
5221 action,
5222 Action::Open {
5223 order_type: OrderType::Market,
5224 ..
5225 } | Action::ClosePosition { .. }
5226 | Action::ClosePartial { .. }
5227 | Action::ScaleIn { .. }
5228 )
5229}
5230
5231fn raw_signal_kind(signal: &RawSignal) -> &'static str {
5232 match signal {
5233 RawSignal::Entry { .. } => "entry",
5234 RawSignal::Close { .. } => "close",
5235 RawSignal::ClosePartial { .. } => "close_partial",
5236 RawSignal::ModifyStoploss { .. } => "modify_stoploss",
5237 RawSignal::MoveStoplossToEntry { .. } => "move_stoploss_to_entry",
5238 RawSignal::AddTarget { .. } => "add_target",
5239 RawSignal::RemoveTarget { .. } => "remove_target",
5240 RawSignal::ModifyTarget { .. } => "modify_target",
5241 RawSignal::AddRule { .. } => "add_rule",
5242 RawSignal::RemoveRule { .. } => "remove_rule",
5243 RawSignal::ScaleIn { .. } => "scale_in",
5244 RawSignal::CancelPending { .. } => "cancel_pending",
5245 RawSignal::CloseAllOf { .. } => "close_all_of",
5246 RawSignal::CloseAll { .. } => "close_all",
5247 RawSignal::CancelAllPending { .. } => "cancel_all_pending",
5248 RawSignal::ModifyAllStoploss { .. } => "modify_all_stoploss",
5249 RawSignal::CloseAllInGroup { .. } => "close_all_in_group",
5250 RawSignal::ModifyAllStoplossInGroup { .. } => "modify_all_stoploss_in_group",
5251 }
5252}
5253
5254#[cfg(test)]
5257mod tests {
5258 use super::*;
5259 use crate::currency::{ConversionRoute, FxPair};
5260 use crate::data_feed::{EventMetadata, FeedEvent, MarketEvent, SeriesRoles, VecFeed};
5261 use crate::profile::{
5262 EntryGeometryPolicy, ManagementProfile, PositionRef, RawSignal, StoplossMode, TargetSource,
5263 };
5264 use chrono::NaiveDate;
5265 use qs_core::types::{CloseReason, FillPurpose, OrderType, Side, TargetSpec};
5266
5267 fn ts(h: u32, m: u32, s: u32) -> chrono::NaiveDateTime {
5268 NaiveDate::from_ymd_opt(2026, 1, 1)
5269 .unwrap()
5270 .and_hms_opt(h, m, s)
5271 .unwrap()
5272 }
5273
5274 fn tick(symbol: &str, bid: f64, ask: f64, time: chrono::NaiveDateTime) -> MarketEvent {
5275 MarketEvent::Tick {
5276 symbol: symbol.into(),
5277 ts: time,
5278 bid,
5279 ask,
5280 }
5281 }
5282
5283 #[test]
5284 fn scheduled_signal_preserves_an_explicit_opaque_action_id() {
5285 let scheduled = ScheduledSignal::new(
5286 7,
5287 ts(10, 0, 0),
5288 ts(10, 0, 1),
5289 RawSignal::CloseAll { ts: ts(10, 0, 0) },
5290 true,
5291 )
5292 .with_action_id("caller-command/opaque:7");
5293
5294 assert_eq!(scheduled.resolved_action_id(), "caller-command/opaque:7");
5295 }
5296
5297 #[test]
5298 fn scheduled_signal_keeps_the_compatible_generated_action_id() {
5299 let scheduled = ScheduledSignal::new(
5300 7,
5301 ts(10, 0, 0),
5302 ts(10, 0, 1),
5303 RawSignal::CloseAll { ts: ts(10, 0, 0) },
5304 false,
5305 );
5306
5307 assert_eq!(scheduled.resolved_action_id(), "signal:00000007");
5308 }
5309
5310 fn test_symbol_spec(symbol: &str) -> qs_symbols::SymbolSpec {
5311 qs_symbols::SymbolSpec {
5312 canonical: symbol.to_ascii_lowercase(),
5313 pip_position: 4,
5314 digits: 5,
5315 category: "forex".into(),
5316 lot_base_units: 100,
5317 lot_step_units: 1,
5318 lot_min_steps: 1,
5319 lot_max_steps: 0,
5320 }
5321 }
5322
5323 fn fixed_lot_config() -> BacktestConfig {
5324 BacktestConfig {
5325 sizing: Some(SizingPolicy::FixedLot { lots: 1.0 }),
5326 symbol_specs: ["EURUSD", "XAUUSD"]
5327 .into_iter()
5328 .map(|symbol| (symbol.to_owned(), test_symbol_spec(symbol)))
5329 .collect(),
5330 ..BacktestConfig::default()
5331 }
5332 }
5333
5334 fn identity_currency_plan(symbol: &str) -> RunCurrencyPlan {
5335 RunCurrencyPlan::new(
5336 "USD",
5337 [symbol.to_owned()].into_iter().collect(),
5338 Default::default(),
5339 [(symbol.to_owned(), "USD".to_owned())]
5340 .into_iter()
5341 .collect(),
5342 [(
5343 "USD".to_owned(),
5344 ConversionRoute::Identity {
5345 currency: "USD".to_owned(),
5346 },
5347 )]
5348 .into_iter()
5349 .collect(),
5350 Vec::new(),
5351 )
5352 .unwrap()
5353 }
5354
5355 struct ScriptedBatchFeed {
5356 batches: VecDeque<Result<Option<TimestampBatch>, &'static str>>,
5357 }
5358
5359 impl FallibleBatchFeed for ScriptedBatchFeed {
5360 type Error = &'static str;
5361
5362 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
5363 self.batches.pop_front().unwrap_or(Ok(None))
5364 }
5365 }
5366
5367 struct CountingBatchFeed {
5368 batches: VecDeque<TimestampBatch>,
5369 polls: std::rc::Rc<std::cell::Cell<usize>>,
5370 }
5371
5372 impl FallibleBatchFeed for CountingBatchFeed {
5373 type Error = Infallible;
5374
5375 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
5376 self.polls.set(self.polls.get() + 1);
5377 Ok(self.batches.pop_front())
5378 }
5379 }
5380
5381 fn primary_batch(event: MarketEvent) -> TimestampBatch {
5382 TimestampBatch {
5383 ts: event.ts(),
5384 events: vec![FeedEvent::new(
5385 event,
5386 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
5387 )],
5388 }
5389 }
5390
5391 fn market_entry(timestamp: NaiveDateTime, symbol: &str, order_type: OrderType) -> RawSignal {
5392 RawSignal::Entry {
5393 ts: timestamp,
5394 symbol: symbol.into(),
5395 side: Side::Buy,
5396 order_type,
5397 price: (order_type == OrderType::Limit).then_some(1.0),
5398 risk_multiplier: 1.0,
5399 stoploss: None,
5400 targets: Vec::new(),
5401 group: None,
5402 trade_id: Some(format!("{symbol}-blocker")),
5403 entry_class: None,
5404 }
5405 }
5406
5407 #[test]
5408 fn future_streaming_matches_materialized_and_stops_without_draining() {
5409 let events = vec![
5410 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5411 tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
5412 tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
5413 ];
5414 let signals = vec![
5415 market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market),
5416 RawSignal::CloseAll { ts: ts(10, 0, 1) },
5417 ];
5418 let config = BacktestConfig {
5419 close_on_finish: false,
5420 ..fixed_lot_config()
5421 };
5422 let mut materialized_feed = VecFeed::new(events.clone());
5423 let materialized = BacktestRunner::new_future(
5424 config.clone(),
5425 FutureQuoteConfig {
5426 mtm_output: MtmOutputPolicy::Full,
5427 ..FutureQuoteConfig::default()
5428 },
5429 )
5430 .run_raw_signals_future(&mut materialized_feed, signals.clone(), None);
5431
5432 let mut stream = ScriptedBatchFeed {
5433 batches: VecDeque::from([
5434 Ok(Some(primary_batch(events[0].clone()))),
5435 Ok(Some(primary_batch(events[1].clone()))),
5436 Ok(Some(primary_batch(events[2].clone()))),
5437 Err("must not drain"),
5438 ]),
5439 };
5440 let mut progress = Vec::new();
5441 let streamed = BacktestRunner::new_future(
5442 config,
5443 FutureQuoteConfig {
5444 mtm_output: MtmOutputPolicy::Full,
5445 ..FutureQuoteConfig::default()
5446 },
5447 )
5448 .run_raw_signals_future_streaming_controlled(
5449 &mut stream,
5450 Some(ts(10, 0, 2)),
5451 signals,
5452 None,
5453 || false,
5454 |update| progress.push(update),
5455 )
5456 .unwrap();
5457
5458 assert_eq!(
5459 serde_json::to_value(&streamed).unwrap(),
5460 serde_json::to_value(&materialized).unwrap()
5461 );
5462 assert_eq!(
5463 stream.batches.len(),
5464 2,
5465 "quiescence must leave the tail unread"
5466 );
5467 assert_eq!(progress.first().unwrap().total_events, 0);
5468 assert_eq!(progress.last().unwrap().processed_events, 2);
5469 assert_eq!(progress.last().unwrap().total_events, 2);
5470 assert_eq!(
5471 streamed.mtm_equity_curve.last().unwrap().ts,
5472 ts(10, 0, 1),
5473 "terminal observation must use the last processed primary timestamp"
5474 );
5475 assert_eq!(
5476 streamed
5477 .mtm_equity_curve
5478 .last()
5479 .unwrap()
5480 .observation_kind
5481 .as_deref(),
5482 Some(EquityObservationKind::QuiescentTermination.as_str())
5483 );
5484 assert_eq!(
5485 streamed
5486 .execution_metadata
5487 .as_ref()
5488 .unwrap()
5489 .tags
5490 .get("termination_reason")
5491 .map(String::as_str),
5492 Some("quiescent")
5493 );
5494 }
5495
5496 #[test]
5497 fn exact_time_close_waits_for_later_symbol_pending_fill() {
5498 let open_ts = ts(10, 0, 0);
5499 let execution_ts = ts(10, 0, 1);
5500 let events = vec![
5501 FeedEvent::new(
5502 tick("XAUUSD", 101.0, 101.0, open_ts),
5503 EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
5504 ),
5505 FeedEvent::new(
5506 tick("EURUSD", 1.1, 1.1, execution_ts),
5507 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
5508 ),
5509 FeedEvent::new(
5510 tick("XAUUSD", 100.0, 100.0, execution_ts),
5511 EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
5512 ),
5513 ];
5514 let signals = vec![
5515 RawSignal::Entry {
5516 ts: open_ts,
5517 symbol: "XAUUSD".into(),
5518 side: Side::Buy,
5519 order_type: OrderType::Limit,
5520 price: Some(100.0),
5521 risk_multiplier: 1.0,
5522 stoploss: None,
5523 targets: Vec::new(),
5524 group: None,
5525 trade_id: Some("later-pending".into()),
5526 entry_class: None,
5527 },
5528 RawSignal::Close {
5529 ts: execution_ts,
5530 position: PositionRef::ByTradeId {
5531 trade_id: "later-pending".into(),
5532 },
5533 },
5534 ];
5535 let mut feed = VecFeed::from_feed_events(events);
5536 let result = BacktestRunner::new_future(
5537 BacktestConfig {
5538 close_on_finish: false,
5539 ..fixed_lot_config()
5540 },
5541 FutureQuoteConfig::default(),
5542 )
5543 .run_raw_signals_future(&mut feed, signals, None);
5544
5545 assert_eq!(
5546 result
5547 .recorded_fills
5548 .iter()
5549 .map(|fill| fill.fill.purpose)
5550 .collect::<Vec<_>>(),
5551 vec![FillPurpose::LimitEntry, FillPurpose::MarketExit]
5552 );
5553 assert_eq!(result.close_events.len(), 1);
5554 assert_eq!(result.close_events[0].reason, CloseReason::Manual);
5555 assert!(result.open_position_snapshots.is_empty());
5556 assert!(result.pending_order_snapshots.is_empty());
5557 }
5558
5559 #[test]
5560 fn exact_time_close_cannot_beat_later_symbol_stoploss() {
5561 let open_ts = ts(10, 0, 0);
5562 let execution_ts = ts(10, 0, 1);
5563 let events = vec![
5564 FeedEvent::new(
5565 tick("XAUUSD", 100.0, 100.0, open_ts),
5566 EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
5567 ),
5568 FeedEvent::new(
5569 tick("EURUSD", 1.1, 1.1, execution_ts),
5570 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
5571 ),
5572 FeedEvent::new(
5573 tick("XAUUSD", 98.0, 98.0, execution_ts),
5574 EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
5575 ),
5576 ];
5577 let signals = vec![
5578 RawSignal::Entry {
5579 ts: open_ts,
5580 symbol: "XAUUSD".into(),
5581 side: Side::Buy,
5582 order_type: OrderType::Market,
5583 price: None,
5584 risk_multiplier: 1.0,
5585 stoploss: Some(99.0),
5586 targets: Vec::new(),
5587 group: None,
5588 trade_id: Some("later-stop".into()),
5589 entry_class: None,
5590 },
5591 RawSignal::Close {
5592 ts: execution_ts,
5593 position: PositionRef::ByTradeId {
5594 trade_id: "later-stop".into(),
5595 },
5596 },
5597 ];
5598 let mut feed = VecFeed::from_feed_events(events);
5599 let result = BacktestRunner::new_future(
5600 BacktestConfig {
5601 close_on_finish: false,
5602 ..fixed_lot_config()
5603 },
5604 FutureQuoteConfig::default(),
5605 )
5606 .run_raw_signals_future(&mut feed, signals, None);
5607
5608 assert_eq!(result.close_events.len(), 1);
5609 assert_eq!(result.close_events[0].reason, CloseReason::Stoploss);
5610 assert_eq!(
5611 result.recorded_fills.last().unwrap().fill.purpose,
5612 FillPurpose::StopLoss
5613 );
5614 assert!(!result.action_dispositions.iter().any(|disposition| {
5615 disposition.action_id.starts_with("signal:00000001")
5616 && disposition.status == crate::ledger::ActionDispositionStatus::Applied
5617 }));
5618 }
5619
5620 #[test]
5621 fn exact_time_multisymbol_closes_preserve_signal_order() {
5622 let open_ts = ts(10, 0, 0);
5623 let close_ts = ts(10, 0, 1);
5624 let mut events = Vec::new();
5625 for (timestamp, row) in [(open_ts, 0), (close_ts, 1)] {
5626 events.push(FeedEvent::new(
5627 tick("EURUSD", 1.1, 1.1, timestamp),
5628 EventMetadata::new(SeriesRoles::PRIMARY, 0, row),
5629 ));
5630 events.push(FeedEvent::new(
5631 tick("XAUUSD", 100.0, 100.0, timestamp),
5632 EventMetadata::new(SeriesRoles::PRIMARY, 1, row),
5633 ));
5634 }
5635 let entry = |symbol: &str, trade_id: &str| RawSignal::Entry {
5636 ts: open_ts,
5637 symbol: symbol.into(),
5638 side: Side::Buy,
5639 order_type: OrderType::Market,
5640 price: None,
5641 risk_multiplier: 1.0,
5642 stoploss: None,
5643 targets: Vec::new(),
5644 group: None,
5645 trade_id: Some(trade_id.into()),
5646 entry_class: None,
5647 };
5648 let close = |trade_id: &str| RawSignal::Close {
5649 ts: close_ts,
5650 position: PositionRef::ByTradeId {
5651 trade_id: trade_id.into(),
5652 },
5653 };
5654 let signals = vec![
5655 entry("XAUUSD", "close-first"),
5656 entry("EURUSD", "close-second"),
5657 close("close-first"),
5658 close("close-second"),
5659 ];
5660 let mut feed = VecFeed::from_feed_events(events);
5661 let result = BacktestRunner::new_future(
5662 BacktestConfig {
5663 close_on_finish: false,
5664 ..fixed_lot_config()
5665 },
5666 FutureQuoteConfig::default(),
5667 )
5668 .run_raw_signals_future(&mut feed, signals, None);
5669
5670 assert_eq!(
5671 result
5672 .close_events
5673 .iter()
5674 .map(|event| event.symbol.as_str())
5675 .collect::<Vec<_>>(),
5676 vec!["XAUUSD", "EURUSD"]
5677 );
5678 assert!(
5679 result
5680 .close_events
5681 .iter()
5682 .all(|event| event.reason == CloseReason::Manual)
5683 );
5684 }
5685
5686 #[test]
5687 fn future_streaming_quiescence_waits_for_all_blockers() {
5688 let run = |events: Vec<MarketEvent>, signals: Vec<RawSignal>, config: BacktestConfig| {
5689 let polls = std::rc::Rc::new(std::cell::Cell::new(0));
5690 let primary_eod = events.last().map(MarketEvent::ts);
5691 let mut feed = CountingBatchFeed {
5692 batches: events.into_iter().map(primary_batch).collect(),
5693 polls: polls.clone(),
5694 };
5695 BacktestRunner::new_future(config, FutureQuoteConfig::default())
5696 .run_raw_signals_future_streaming_controlled(
5697 &mut feed,
5698 primary_eod,
5699 signals,
5700 None,
5701 || false,
5702 |_| {},
5703 )
5704 .unwrap();
5705 polls.get()
5706 };
5707 let eur_events = vec![
5708 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5709 tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
5710 tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
5711 ];
5712
5713 let immediately_quiescent = run(
5714 eur_events.clone(),
5715 vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
5716 BacktestConfig::default(),
5717 );
5718 assert_eq!(immediately_quiescent, 1);
5719
5720 let scheduled = run(
5721 eur_events.clone(),
5722 vec![RawSignal::CloseAll { ts: ts(10, 0, 2) }],
5723 BacktestConfig::default(),
5724 );
5725 assert_eq!(scheduled, 3, "scheduled signals must block termination");
5726
5727 let mut two_symbol_config = BacktestConfig {
5728 close_on_finish: false,
5729 ..fixed_lot_config()
5730 };
5731 two_symbol_config
5732 .symbol_specs
5733 .insert("GBPUSD".into(), test_symbol_spec("GBPUSD"));
5734 let queued = run(
5735 vec![
5736 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5737 tick("GBPUSD", 1.2500, 1.2502, ts(10, 0, 1)),
5738 tick("GBPUSD", 1.2501, 1.2503, ts(10, 0, 2)),
5739 ],
5740 vec![
5741 market_entry(ts(10, 0, 0), "GBPUSD", OrderType::Market),
5742 RawSignal::CloseAll { ts: ts(10, 0, 1) },
5743 ],
5744 two_symbol_config,
5745 );
5746 assert_eq!(
5747 queued, 2,
5748 "queued actions must wait for an eligible symbol quote"
5749 );
5750
5751 let open = run(
5752 eur_events.clone(),
5753 vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market)],
5754 BacktestConfig {
5755 close_on_finish: false,
5756 ..fixed_lot_config()
5757 },
5758 );
5759 assert_eq!(
5760 open, 4,
5761 "open positions must consume the stream through EOD"
5762 );
5763
5764 let pending = run(
5765 eur_events,
5766 vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Limit)],
5767 fixed_lot_config(),
5768 );
5769 assert_eq!(
5770 pending, 4,
5771 "pending orders must consume the stream through EOD"
5772 );
5773 }
5774
5775 #[test]
5776 fn future_mtm_output_policies_bound_curve_and_validate_before_feed_use() {
5777 assert_eq!(
5778 FutureQuoteConfig::default().mtm_output,
5779 MtmOutputPolicy::Bounded { max_points: 4_096 }
5780 );
5781 let events: Vec<_> = (0..12)
5782 .map(|second| tick("EURUSD", 100.0, 100.0, ts(10, 0, second)))
5783 .collect();
5784
5785 let pending = RawSignal::Entry {
5786 ts: ts(10, 0, 0),
5787 symbol: "EURUSD".into(),
5788 side: Side::Buy,
5789 order_type: OrderType::Limit,
5790 price: Some(90.0),
5791 risk_multiplier: 1.0,
5792 stoploss: None,
5793 targets: Vec::new(),
5794 group: None,
5795 trade_id: Some("mtm-policy-blocker".into()),
5796 entry_class: None,
5797 };
5798 let run = |policy| {
5799 let mut feed = VecFeed::new(events.clone());
5800 BacktestRunner::new_future(
5801 fixed_lot_config(),
5802 FutureQuoteConfig {
5803 mtm_output: policy,
5804 ..FutureQuoteConfig::default()
5805 },
5806 )
5807 .run_raw_signals_future(&mut feed, vec![pending.clone()], None)
5808 };
5809
5810 let none = run(MtmOutputPolicy::None);
5811 assert!(none.mtm_equity_curve.is_empty());
5812 assert_eq!(none.mtm_output_summary.observed_points, 13);
5813 assert_eq!(none.mtm_output_summary.omitted_points, 13);
5814
5815 let bounded = run(MtmOutputPolicy::Bounded { max_points: 8 });
5816 assert_eq!(bounded.mtm_equity_curve.len(), 8);
5817 assert_eq!(bounded.mtm_output_summary.observed_points, 13);
5818 assert_eq!(bounded.mtm_output_summary.retained_points, 8);
5819 assert_eq!(bounded.mtm_output_summary.omitted_points, 5);
5820
5821 let full = run(MtmOutputPolicy::Full);
5822 assert_eq!(full.mtm_equity_curve.len(), 13);
5823 assert_eq!(full.mtm_output_summary.observed_points, 13);
5824 assert_eq!(full.mtm_output_summary.omitted_points, 0);
5825 assert_eq!(
5826 full.mtm_equity_curve
5827 .iter()
5828 .filter(|point| {
5829 point.observation_kind.as_deref()
5830 == Some(EquityObservationKind::PostOutput.as_str())
5831 })
5832 .count(),
5833 0
5834 );
5835
5836 let mut invalid_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5837 let rejected = BacktestRunner::new_future(
5838 BacktestConfig::default(),
5839 FutureQuoteConfig {
5840 mtm_output: MtmOutputPolicy::Bounded { max_points: 7 },
5841 ..FutureQuoteConfig::default()
5842 },
5843 )
5844 .run_raw_signals_future(&mut invalid_feed, Vec::new(), None);
5845 assert_eq!(invalid_feed.remaining(), 1);
5846 assert!(rejected.action_dispositions.iter().any(|disposition| {
5847 disposition.action_id == "configuration"
5848 && disposition
5849 .reason
5850 .as_deref()
5851 .is_some_and(|reason| reason.contains("MTM max_points"))
5852 }));
5853 }
5854
5855 #[test]
5856 fn future_mtm_records_changed_post_output_observation_kind() {
5857 let mut feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5858 let signal = RawSignal::Entry {
5859 ts: ts(10, 0, 0),
5860 symbol: "EURUSD".into(),
5861 side: Side::Buy,
5862 order_type: OrderType::Market,
5863 price: None,
5864 risk_multiplier: 1.0,
5865 stoploss: None,
5866 targets: Vec::new(),
5867 group: None,
5868 trade_id: Some("mtm-kind".into()),
5869 entry_class: None,
5870 };
5871 let result = BacktestRunner::new_future(
5872 BacktestConfig {
5873 close_on_finish: false,
5874 ..fixed_lot_config()
5875 },
5876 FutureQuoteConfig {
5877 mtm_output: MtmOutputPolicy::Full,
5878 ..FutureQuoteConfig::default()
5879 },
5880 )
5881 .run_raw_signals_future(&mut feed, vec![signal], None);
5882
5883 let kinds: Vec<_> = result
5884 .mtm_equity_curve
5885 .iter()
5886 .filter_map(|point| point.observation_kind.as_deref())
5887 .collect();
5888 assert_eq!(
5889 kinds,
5890 vec![
5891 EquityObservationKind::PreSettlement.as_str(),
5892 EquityObservationKind::PostOutput.as_str(),
5893 EquityObservationKind::EndOfData.as_str(),
5894 ]
5895 );
5896 assert_eq!(
5897 result
5898 .execution_metadata
5899 .as_ref()
5900 .unwrap()
5901 .tags
5902 .get("termination_reason")
5903 .map(String::as_str),
5904 Some("end_of_data")
5905 );
5906 }
5907
5908 #[test]
5909 fn future_fallible_batch_feed_propagates_source_error() {
5910 let batch = TimestampBatch {
5911 ts: ts(10, 0, 0),
5912 events: vec![FeedEvent::new(
5913 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
5914 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
5915 )],
5916 };
5917 let mut feed = ScriptedBatchFeed {
5918 batches: VecDeque::from([Ok(Some(batch)), Err("feed failed")]),
5919 };
5920 let result =
5921 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
5922 .run_raw_signals_future_fallible(&mut feed, Vec::new(), None);
5923
5924 assert!(matches!(result, Err("feed failed")));
5925 }
5926
5927 struct BuyOnceStrategy {
5931 entered: bool,
5932 }
5933
5934 impl BuyOnceStrategy {
5935 fn new() -> Self {
5936 Self { entered: false }
5937 }
5938 }
5939
5940 impl Strategy for BuyOnceStrategy {
5941 fn on_event(&mut self, event: &MarketEvent) -> Vec<Action> {
5942 if self.entered {
5943 return vec![];
5944 }
5945 if let MarketEvent::Tick { symbol, ask, .. } = event {
5946 self.entered = true;
5947 vec![Action::Open {
5948 symbol: symbol.clone(),
5949 side: Side::Buy,
5950 order_type: OrderType::Market,
5951 price: Some(*ask),
5952 size: 1.0,
5953 stoploss: Some(*ask - 0.0050),
5954 targets: vec![TargetSpec {
5955 price: *ask + 0.0050,
5956 close_ratio: 1.0,
5957 }],
5958 rules: vec![],
5959 group: None,
5960 trade_id: None,
5961 }]
5962 } else {
5963 vec![]
5964 }
5965 }
5966
5967 fn on_finished(&mut self) -> Vec<Action> {
5968 vec![]
5971 }
5972 }
5973
5974 #[test]
5977 fn strategy_backtest_tp_hit() {
5978 let events = vec![
5979 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
5980 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
5981 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
5982 tick("EURUSD", 1.0890, 1.0892, ts(10, 0, 3)),
5983 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 4)),
5985 ];
5986 let mut feed = VecFeed::new(events);
5987 let mut strategy = BuyOnceStrategy::new();
5988
5989 let config = BacktestConfig {
5990 initial_balance: 10_000.0,
5991 close_on_finish: true,
5992 ..Default::default()
5993 };
5994 let runner = BacktestRunner::new(config);
5995 let result = runner.run_strategy(&mut feed, &mut strategy);
5996
5997 assert_eq!(result.total_trades, 1);
5998 assert_eq!(result.winning_trades, 1);
5999 assert!(result.total_pnl > 0.0);
6000 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6001 }
6002
6003 #[test]
6004 fn strategy_backtest_sl_hit() {
6005 let events = vec![
6006 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6007 tick("EURUSD", 1.0830, 1.0832, ts(10, 0, 1)),
6008 tick("EURUSD", 1.0799, 1.0801, ts(10, 0, 2)),
6010 ];
6011 let mut feed = VecFeed::new(events);
6012 let mut strategy = BuyOnceStrategy::new();
6013
6014 let runner = BacktestRunner::with_defaults();
6015 let result = runner.run_strategy(&mut feed, &mut strategy);
6016
6017 assert_eq!(result.total_trades, 1);
6018 assert_eq!(result.losing_trades, 1);
6019 assert!(result.total_pnl < 0.0);
6020 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6021 }
6022
6023 #[test]
6024 fn strategy_close_on_finish() {
6025 let events = vec![
6027 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6028 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6029 tick("EURUSD", 1.0852, 1.0854, ts(10, 0, 2)),
6030 ];
6031 let mut feed = VecFeed::new(events);
6032 let mut strategy = BuyOnceStrategy::new();
6033
6034 let config = BacktestConfig {
6035 initial_balance: 10_000.0,
6036 close_on_finish: true,
6037 ..Default::default()
6038 };
6039 let runner = BacktestRunner::new(config);
6040 let result = runner.run_strategy(&mut feed, &mut strategy);
6041
6042 assert_eq!(result.total_trades, 1);
6043 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6044 }
6045
6046 #[test]
6047 fn strategy_no_close_on_finish() {
6048 let events = vec![
6049 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6050 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6051 ];
6052 let mut feed = VecFeed::new(events);
6053 let mut strategy = BuyOnceStrategy::new();
6054
6055 let config = BacktestConfig {
6056 initial_balance: 10_000.0,
6057 close_on_finish: false,
6058 ..Default::default()
6059 };
6060 let runner = BacktestRunner::new(config);
6061 let result = runner.run_strategy(&mut feed, &mut strategy);
6062
6063 assert_eq!(result.total_trades, 0);
6065 }
6066
6067 #[test]
6070 fn legacy_unprofiled_targets_default_to_equal_weights() {
6071 let events = vec![
6072 tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
6073 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
6074 tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
6075 ];
6076 let mut feed = VecFeed::new(events);
6077 let signals = vec![RawSignal::Entry {
6078 ts: ts(10, 0, 0),
6079 symbol: "EURUSD".into(),
6080 side: Side::Buy,
6081 order_type: OrderType::Market,
6082 price: Some(1.0000),
6083 risk_multiplier: 1.0,
6084 stoploss: None,
6085 targets: vec![1.1000, 1.2000],
6086 group: None,
6087 trade_id: Some("equal-targets".into()),
6088 entry_class: None,
6089 }];
6090
6091 let result = BacktestRunner::new(BacktestConfig {
6092 close_on_finish: false,
6093 ..fixed_lot_config()
6094 })
6095 .run_raw_signals(&mut feed, signals, None);
6096
6097 assert_eq!(result.trade_log.len(), 2);
6098 assert!(
6099 result
6100 .trade_log
6101 .iter()
6102 .all(|trade| (trade.size - 0.5).abs() < f64::EPSILON)
6103 );
6104 assert!(
6105 result
6106 .trade_log
6107 .iter()
6108 .all(|trade| trade.close_reason == CloseReason::Target)
6109 );
6110 }
6111
6112 #[test]
6113 fn legacy_atomic_target_modification_retains_profile_ratio() {
6114 let events = vec![
6115 tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
6116 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
6117 tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
6118 tick("EURUSD", 1.3000, 1.3000, ts(10, 0, 3)),
6119 ];
6120 let mut feed = VecFeed::new(events);
6121 let position = PositionRef::ByTradeId {
6122 trade_id: "modified-target".into(),
6123 };
6124 let signals = vec![
6125 RawSignal::Entry {
6126 ts: ts(10, 0, 0),
6127 symbol: "EURUSD".into(),
6128 side: Side::Buy,
6129 order_type: OrderType::Market,
6130 price: Some(1.0000),
6131 risk_multiplier: 1.0,
6132 stoploss: None,
6133 targets: vec![1.1000, 1.3000],
6134 group: None,
6135 trade_id: Some("modified-target".into()),
6136 entry_class: None,
6137 },
6138 RawSignal::ModifyTarget {
6139 ts: ts(10, 0, 0),
6140 position,
6141 old_price: 1.1000,
6142 new_price: 1.2000,
6143 },
6144 ];
6145 let profile = ManagementProfile {
6146 name: "non-default-ratios".into(),
6147 target_selection: None,
6148 use_targets: vec![1, 2],
6149 close_ratios: vec![0.25, 0.75],
6150 target_source: TargetSource::FromSignal,
6151 stoploss_mode: StoplossMode::FromSignal,
6152 rules: vec![],
6153 group_override: None,
6154 let_remainder_run: false,
6155 entry_geometry: EntryGeometryPolicy::Strict,
6156 };
6157
6158 let result = BacktestRunner::new(BacktestConfig {
6159 close_on_finish: false,
6160 ..fixed_lot_config()
6161 })
6162 .run_raw_signals(&mut feed, signals, Some(&profile));
6163
6164 assert_eq!(result.trade_log.len(), 2);
6165 assert!((result.trade_log[0].exit_price - 1.2000).abs() < f64::EPSILON);
6166 assert!((result.trade_log[0].size - 0.25).abs() < f64::EPSILON);
6167 assert!((result.trade_log[1].exit_price - 1.3000).abs() < f64::EPSILON);
6168 assert!((result.trade_log[1].size - 0.75).abs() < f64::EPSILON);
6169 }
6170
6171 #[test]
6172 fn run_raw_signals_entry_only() {
6173 let events = vec![
6174 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6175 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6176 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6177 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
6178 ];
6179 let mut feed = VecFeed::new(events);
6180
6181 let raw_signals = vec![RawSignal::Entry {
6182 ts: ts(10, 0, 0),
6183 symbol: "EURUSD".into(),
6184 side: Side::Buy,
6185 order_type: OrderType::Market,
6186 price: Some(1.0850),
6187 risk_multiplier: 1.0,
6188 stoploss: Some(1.0800),
6189 targets: vec![1.0900],
6190 group: None,
6191 trade_id: None,
6192 entry_class: None,
6193 }];
6194
6195 let runner = BacktestRunner::new(fixed_lot_config());
6196 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6197
6198 assert_eq!(result.total_trades, 1);
6199 assert_eq!(result.winning_trades, 1);
6200 }
6201
6202 #[test]
6203 fn run_raw_signals_open_then_close() {
6204 let events = vec![
6205 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6206 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6207 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6208 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6209 ];
6210 let mut feed = VecFeed::new(events);
6211
6212 let raw_signals = vec![
6213 RawSignal::Entry {
6214 ts: ts(10, 0, 0),
6215 symbol: "EURUSD".into(),
6216 side: Side::Buy,
6217 order_type: OrderType::Market,
6218 price: Some(1.0850),
6219 risk_multiplier: 1.0,
6220 stoploss: None,
6221 targets: vec![],
6222 group: None,
6223 trade_id: Some("t1".into()),
6224 entry_class: None,
6225 },
6226 RawSignal::Close {
6227 ts: ts(10, 0, 2),
6228 position: PositionRef::ByTradeId {
6229 trade_id: "t1".into(),
6230 },
6231 },
6232 ];
6233
6234 let config = BacktestConfig {
6235 initial_balance: 10_000.0,
6236 close_on_finish: false,
6237 ..fixed_lot_config()
6238 };
6239 let runner = BacktestRunner::new(config);
6240 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6241
6242 assert_eq!(result.total_trades, 1);
6243 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6244 }
6245
6246 #[test]
6247 fn run_raw_signals_open_then_modify_sl() {
6248 let events = vec![
6250 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6251 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
6252 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 2)),
6254 tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
6256 ];
6257 let mut feed = VecFeed::new(events);
6258
6259 let raw_signals = vec![
6260 RawSignal::Entry {
6261 ts: ts(10, 0, 0),
6262 symbol: "EURUSD".into(),
6263 side: Side::Buy,
6264 order_type: OrderType::Market,
6265 price: Some(1.0850),
6266 risk_multiplier: 1.0,
6267 stoploss: Some(1.0800),
6268 targets: vec![],
6269 group: None,
6270 trade_id: Some("t1".into()),
6271 entry_class: None,
6272 },
6273 RawSignal::ModifyStoploss {
6274 ts: ts(10, 0, 2),
6275 position: PositionRef::ByTradeId {
6276 trade_id: "t1".into(),
6277 },
6278 price: 1.0840,
6279 },
6280 ];
6281
6282 let config = BacktestConfig {
6283 initial_balance: 10_000.0,
6284 close_on_finish: true,
6285 ..fixed_lot_config()
6286 };
6287 let runner = BacktestRunner::new(config);
6288 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6289
6290 assert_eq!(result.total_trades, 1);
6291 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6292 }
6293
6294 #[test]
6295 fn run_raw_signals_open_then_partial_close() {
6296 let events = vec![
6297 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6298 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
6299 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6300 tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 3)),
6301 ];
6302 let mut feed = VecFeed::new(events);
6303
6304 let raw_signals = vec![
6305 RawSignal::Entry {
6306 ts: ts(10, 0, 0),
6307 symbol: "EURUSD".into(),
6308 side: Side::Buy,
6309 order_type: OrderType::Market,
6310 price: Some(1.0850),
6311 risk_multiplier: 1.0,
6312 stoploss: None,
6313 targets: vec![],
6314 group: None,
6315 trade_id: Some("t1".into()),
6316 entry_class: None,
6317 },
6318 RawSignal::ClosePartial {
6319 ts: ts(10, 0, 1),
6320 position: PositionRef::ByTradeId {
6321 trade_id: "t1".into(),
6322 },
6323 ratio: 0.5,
6324 },
6325 ];
6326
6327 let config = BacktestConfig {
6328 initial_balance: 10_000.0,
6329 close_on_finish: true,
6330 ..fixed_lot_config()
6331 };
6332 let runner = BacktestRunner::new(config);
6333 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6334
6335 assert!(result.total_trades >= 1);
6337 }
6338
6339 #[test]
6340 fn run_raw_signals_group_workflow() {
6341 let events = vec![
6342 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6343 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6344 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6345 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6346 tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 4)),
6347 ];
6348 let mut feed = VecFeed::new(events);
6349
6350 let raw_signals = vec![
6351 RawSignal::Entry {
6353 ts: ts(10, 0, 0),
6354 symbol: "EURUSD".into(),
6355 side: Side::Buy,
6356 order_type: OrderType::Market,
6357 price: Some(1.0850),
6358 risk_multiplier: 1.0,
6359 stoploss: None,
6360 targets: vec![],
6361 group: Some("grp1".into()),
6362 trade_id: Some("t1".into()),
6363 entry_class: None,
6364 },
6365 RawSignal::Entry {
6366 ts: ts(10, 0, 1),
6367 symbol: "EURUSD".into(),
6368 side: Side::Buy,
6369 order_type: OrderType::Market,
6370 price: Some(1.0857),
6371 risk_multiplier: 1.0,
6372 stoploss: None,
6373 targets: vec![],
6374 group: Some("grp1".into()),
6375 trade_id: Some("t2".into()),
6376 entry_class: None,
6377 },
6378 RawSignal::CloseAllInGroup {
6380 ts: ts(10, 0, 3),
6381 group_id: "grp1".into(),
6382 },
6383 ];
6384
6385 let config = BacktestConfig {
6386 initial_balance: 10_000.0,
6387 close_on_finish: false,
6388 ..fixed_lot_config()
6389 };
6390 let runner = BacktestRunner::new(config);
6391 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6392
6393 assert_eq!(result.total_trades, 2);
6394 }
6395
6396 #[test]
6397 fn run_raw_signals_close_all_of_symbol() {
6398 let events = vec![
6399 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6400 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6401 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6402 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6403 ];
6404 let mut feed = VecFeed::new(events);
6405
6406 let raw_signals = vec![
6407 RawSignal::Entry {
6408 ts: ts(10, 0, 0),
6409 symbol: "EURUSD".into(),
6410 side: Side::Buy,
6411 order_type: OrderType::Market,
6412 price: Some(1.0850),
6413 risk_multiplier: 1.0,
6414 stoploss: None,
6415 targets: vec![],
6416 group: None,
6417 trade_id: Some("t1".into()),
6418 entry_class: None,
6419 },
6420 RawSignal::Entry {
6421 ts: ts(10, 0, 0),
6422 symbol: "EURUSD".into(),
6423 side: Side::Buy,
6424 order_type: OrderType::Market,
6425 price: Some(1.0850),
6426 risk_multiplier: 0.5,
6427 stoploss: None,
6428 targets: vec![],
6429 group: None,
6430 trade_id: Some("t2".into()),
6431 entry_class: None,
6432 },
6433 RawSignal::CloseAllOf {
6434 ts: ts(10, 0, 2),
6435 symbol: "EURUSD".into(),
6436 },
6437 ];
6438
6439 let config = BacktestConfig {
6440 initial_balance: 10_000.0,
6441 close_on_finish: false,
6442 ..fixed_lot_config()
6443 };
6444 let runner = BacktestRunner::new(config);
6445 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6446
6447 assert_eq!(result.total_trades, 2);
6448 }
6449
6450 #[test]
6451 fn run_raw_signals_with_profile() {
6452 let events = vec![
6453 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6454 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6455 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
6456 ];
6457 let mut feed = VecFeed::new(events);
6458
6459 let profile = ManagementProfile {
6460 name: "test".into(),
6461 target_selection: None,
6462 use_targets: vec![1],
6463 close_ratios: vec![1.0],
6464 target_source: TargetSource::FromSignal,
6465 stoploss_mode: StoplossMode::FromSignal,
6466 rules: vec![],
6467 group_override: None,
6468 let_remainder_run: false,
6469 entry_geometry: EntryGeometryPolicy::Strict,
6470 };
6471
6472 let raw_signals = vec![RawSignal::Entry {
6473 ts: ts(10, 0, 0),
6474 symbol: "EURUSD".into(),
6475 side: Side::Buy,
6476 order_type: OrderType::Market,
6477 price: Some(1.0850),
6478 risk_multiplier: 1.0,
6479 stoploss: Some(1.0800),
6480 targets: vec![1.0900],
6481 group: None,
6482 trade_id: Some("t1".into()),
6483 entry_class: None,
6484 }];
6485
6486 let runner = BacktestRunner::new(fixed_lot_config());
6487 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6488
6489 assert_eq!(result.total_trades, 1);
6490 assert_eq!(result.winning_trades, 1);
6491 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6492 }
6493
6494 #[test]
6495 fn run_raw_signals_with_profile_preserves_trade_id() {
6496 let events = vec![
6499 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6500 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6501 ];
6502 let mut feed = VecFeed::new(events);
6503
6504 let profile = ManagementProfile {
6505 name: "test".into(),
6506 target_selection: None,
6507 use_targets: vec![1],
6508 close_ratios: vec![1.0],
6509 target_source: TargetSource::FromSignal,
6510 stoploss_mode: StoplossMode::FromSignal,
6511 rules: vec![],
6512 group_override: None,
6513 let_remainder_run: false,
6514 entry_geometry: EntryGeometryPolicy::Strict,
6515 };
6516
6517 let raw_signals = vec![
6518 RawSignal::Entry {
6519 ts: ts(10, 0, 0),
6520 symbol: "EURUSD".into(),
6521 side: Side::Buy,
6522 order_type: OrderType::Market,
6523 price: Some(1.0850),
6524 risk_multiplier: 1.0,
6525 stoploss: Some(1.0800),
6526 targets: vec![1.0900],
6527 group: None,
6528 trade_id: Some("msg-100".into()),
6529 entry_class: None,
6530 },
6531 RawSignal::Close {
6532 ts: ts(10, 0, 1),
6533 position: PositionRef::ByTradeId {
6534 trade_id: "msg-100".into(),
6535 },
6536 },
6537 ];
6538
6539 let runner = BacktestRunner::new(fixed_lot_config());
6540 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6541
6542 assert_eq!(result.total_trades, 1);
6543 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6544 }
6545
6546 #[test]
6547 fn run_raw_signals_no_profile() {
6548 let events = vec![
6550 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6551 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6552 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6553 ];
6554 let mut feed = VecFeed::new(events);
6555
6556 let raw_signals = vec![RawSignal::Entry {
6557 ts: ts(10, 0, 0),
6558 symbol: "EURUSD".into(),
6559 side: Side::Buy,
6560 order_type: OrderType::Market,
6561 price: Some(1.0850),
6562 risk_multiplier: 1.0,
6563 stoploss: None,
6564 targets: vec![],
6565 group: None,
6566 trade_id: None,
6567 entry_class: None,
6568 }];
6569
6570 let config = BacktestConfig {
6571 initial_balance: 10_000.0,
6572 close_on_finish: true,
6573 ..fixed_lot_config()
6574 };
6575 let runner = BacktestRunner::new(config);
6576 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6577
6578 assert_eq!(result.total_trades, 1);
6579 }
6580
6581 #[test]
6582 fn run_raw_signals_last_on_symbol_resolution() {
6583 let events = vec![
6585 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6586 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6587 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6588 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6589 ];
6590 let mut feed = VecFeed::new(events);
6591
6592 let raw_signals = vec![
6593 RawSignal::Entry {
6594 ts: ts(10, 0, 0),
6595 symbol: "EURUSD".into(),
6596 side: Side::Buy,
6597 order_type: OrderType::Market,
6598 price: Some(1.0850),
6599 risk_multiplier: 1.0,
6600 stoploss: None,
6601 targets: vec![],
6602 group: None,
6603 trade_id: Some("t1".into()),
6604 entry_class: None,
6605 },
6606 RawSignal::Entry {
6607 ts: ts(10, 0, 1),
6608 symbol: "EURUSD".into(),
6609 side: Side::Buy,
6610 order_type: OrderType::Market,
6611 price: Some(1.0857),
6612 risk_multiplier: 1.0,
6613 stoploss: None,
6614 targets: vec![],
6615 group: None,
6616 trade_id: Some("t2".into()),
6617 entry_class: None,
6618 },
6619 RawSignal::Close {
6621 ts: ts(10, 0, 2),
6622 position: PositionRef::ByTradeId {
6623 trade_id: "t2".into(),
6624 },
6625 },
6626 ];
6627
6628 let config = BacktestConfig {
6629 initial_balance: 10_000.0,
6630 close_on_finish: true,
6631 ..fixed_lot_config()
6632 };
6633 let runner = BacktestRunner::new(config);
6634 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6635
6636 assert_eq!(result.total_trades, 2);
6638 }
6639
6640 #[test]
6641 fn run_raw_signals_unresolved_ref_skipped() {
6642 let events = vec![
6644 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6645 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6646 ];
6647 let mut feed = VecFeed::new(events);
6648
6649 let raw_signals = vec![RawSignal::Close {
6650 ts: ts(10, 0, 0),
6651 position: PositionRef::ByTradeId {
6652 trade_id: "nonexistent".into(),
6653 },
6654 }];
6655
6656 let config = BacktestConfig {
6657 initial_balance: 10_000.0,
6658 close_on_finish: false,
6659 ..Default::default()
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, 0);
6666 }
6667
6668 #[test]
6671 fn signal_replay_basic() {
6672 let events = vec![
6673 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6674 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6675 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6676 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
6678 ];
6679 let mut feed = VecFeed::new(events);
6680
6681 let raw_signals = vec![RawSignal::Entry {
6682 ts: ts(10, 0, 0),
6683 symbol: "EURUSD".into(),
6684 side: Side::Buy,
6685 order_type: OrderType::Market,
6686 price: Some(1.0850),
6687 risk_multiplier: 1.0,
6688 stoploss: Some(1.0800),
6689 targets: vec![1.0900],
6690 group: None,
6691 trade_id: None,
6692 entry_class: None,
6693 }];
6694
6695 let runner = BacktestRunner::new(fixed_lot_config());
6696 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6697
6698 assert_eq!(result.total_trades, 1);
6699 assert_eq!(result.winning_trades, 1);
6700 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6701 }
6702
6703 #[test]
6704 fn signal_replay_multiple_signals() {
6705 let events = vec![
6706 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6707 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6708 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
6710 tick("EURUSD", 1.0910, 1.0912, ts(10, 0, 3)),
6711 tick("EURUSD", 1.0920, 1.0922, ts(10, 0, 4)),
6712 ];
6713 let mut feed = VecFeed::new(events);
6714
6715 let raw_signals = vec![
6716 RawSignal::Entry {
6717 ts: ts(10, 0, 0),
6718 symbol: "EURUSD".into(),
6719 side: Side::Buy,
6720 order_type: OrderType::Market,
6721 price: Some(1.0850),
6722 risk_multiplier: 1.0,
6723 stoploss: Some(1.0800),
6724 targets: vec![1.0900],
6725 group: None,
6726 trade_id: Some("t1".into()),
6727 entry_class: None,
6728 },
6729 RawSignal::Entry {
6730 ts: ts(10, 0, 1),
6731 symbol: "EURUSD".into(),
6732 side: Side::Buy,
6733 order_type: OrderType::Market,
6734 price: Some(1.0857),
6735 risk_multiplier: 1.0,
6736 stoploss: None,
6737 targets: vec![],
6738 group: None,
6739 trade_id: Some("t2".into()),
6740 entry_class: None,
6741 },
6742 ];
6743
6744 let config = BacktestConfig {
6745 initial_balance: 10_000.0,
6746 close_on_finish: true,
6747 ..fixed_lot_config()
6748 };
6749 let runner = BacktestRunner::new(config);
6750 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6751
6752 assert!(result.total_trades >= 2);
6754 }
6755
6756 #[test]
6757 fn signal_replay_signal_before_data_filtered() {
6758 let events = vec![
6765 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6766 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6767 ];
6768 let mut feed = VecFeed::new(events);
6769
6770 let raw_signals = vec![RawSignal::Entry {
6771 ts: ts(9, 0, 0), symbol: "EURUSD".into(),
6773 side: Side::Buy,
6774 order_type: OrderType::Market,
6775 price: Some(1.0850),
6776 risk_multiplier: 1.0,
6777 stoploss: None,
6778 targets: vec![1.0900],
6779 group: None,
6780 trade_id: None,
6781 entry_class: None,
6782 }];
6783
6784 let runner = BacktestRunner::new(fixed_lot_config());
6785 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6786
6787 assert_eq!(result.total_trades, 1);
6788 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6789 }
6790
6791 #[test]
6792 fn empty_feed_empty_result() {
6793 let mut feed = VecFeed::new(vec![]);
6794 let mut strategy = BuyOnceStrategy::new();
6795
6796 let runner = BacktestRunner::with_defaults();
6797 let result = runner.run_strategy(&mut feed, &mut strategy);
6798
6799 assert_eq!(result.total_trades, 0);
6800 assert!((result.final_balance - 10_000.0).abs() < f64::EPSILON);
6801 }
6802
6803 #[test]
6804 fn report_display_does_not_panic() {
6805 let events = vec![
6806 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6807 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6808 ];
6809 let mut feed = VecFeed::new(events);
6810 let mut strategy = BuyOnceStrategy::new();
6811
6812 let runner = BacktestRunner::with_defaults();
6813 let result = runner.run_strategy(&mut feed, &mut strategy);
6814
6815 let _display = format!("{result}");
6816 }
6817
6818 #[test]
6819 fn run_raw_signals_with_profile_open_then_modify_sl_by_trade_id() {
6820 let events = vec![
6821 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6822 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6823 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6824 tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
6825 ];
6826 let mut feed = VecFeed::new(events);
6827
6828 let profile = ManagementProfile {
6829 name: "test".into(),
6830 target_selection: None,
6831 use_targets: vec![1],
6832 close_ratios: vec![1.0],
6833 target_source: TargetSource::FromSignal,
6834 stoploss_mode: StoplossMode::FromSignal,
6835 rules: vec![],
6836 group_override: None,
6837 let_remainder_run: false,
6838 entry_geometry: EntryGeometryPolicy::Strict,
6839 };
6840
6841 let raw_signals = vec![
6842 RawSignal::Entry {
6843 ts: ts(10, 0, 0),
6844 symbol: "EURUSD".into(),
6845 side: Side::Buy,
6846 order_type: OrderType::Market,
6847 price: Some(1.0850),
6848 risk_multiplier: 1.0,
6849 stoploss: Some(1.0800),
6850 targets: vec![1.0900],
6851 group: None,
6852 trade_id: Some("t1".into()),
6853 entry_class: None,
6854 },
6855 RawSignal::ModifyStoploss {
6856 ts: ts(10, 0, 2),
6857 position: PositionRef::ByTradeId {
6858 trade_id: "t1".into(),
6859 },
6860 price: 1.0840,
6861 },
6862 ];
6863
6864 let config = BacktestConfig {
6865 initial_balance: 10_000.0,
6866 close_on_finish: false,
6867 ..fixed_lot_config()
6868 };
6869 let runner = BacktestRunner::new(config);
6870 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6871
6872 assert_eq!(result.total_trades, 1);
6873 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6874 }
6875
6876 #[test]
6877 fn run_raw_signals_with_profile_open_then_close_partial_by_trade_id() {
6878 let events = vec![
6879 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6880 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6881 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6882 ];
6883 let mut feed = VecFeed::new(events);
6884
6885 let profile = ManagementProfile {
6886 name: "test".into(),
6887 target_selection: None,
6888 use_targets: vec![1],
6889 close_ratios: vec![1.0],
6890 target_source: TargetSource::FromSignal,
6891 stoploss_mode: StoplossMode::FromSignal,
6892 rules: vec![],
6893 group_override: None,
6894 let_remainder_run: false,
6895 entry_geometry: EntryGeometryPolicy::Strict,
6896 };
6897
6898 let raw_signals = vec![
6899 RawSignal::Entry {
6900 ts: ts(10, 0, 0),
6901 symbol: "EURUSD".into(),
6902 side: Side::Buy,
6903 order_type: OrderType::Market,
6904 price: Some(1.0850),
6905 risk_multiplier: 1.0,
6906 stoploss: Some(1.0800),
6907 targets: vec![1.0900],
6908 group: None,
6909 trade_id: Some("t1".into()),
6910 entry_class: None,
6911 },
6912 RawSignal::ClosePartial {
6913 ts: ts(10, 0, 1),
6914 position: PositionRef::ByTradeId {
6915 trade_id: "t1".into(),
6916 },
6917 ratio: 0.5,
6918 },
6919 ];
6920
6921 let config = BacktestConfig {
6922 initial_balance: 10_000.0,
6923 close_on_finish: false,
6924 ..fixed_lot_config()
6925 };
6926 let runner = BacktestRunner::new(config);
6927 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6928
6929 assert!(result.total_trades >= 1);
6931 }
6932
6933 #[test]
6934 fn run_raw_signals_multi_position_by_trade_id_with_profile() {
6935 let events = vec![
6939 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6940 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6941 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6942 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6943 ];
6944 let mut feed = VecFeed::new(events);
6945
6946 let profile = ManagementProfile {
6947 name: "test".into(),
6948 target_selection: None,
6949 use_targets: vec![1],
6950 close_ratios: vec![1.0],
6951 target_source: TargetSource::FromSignal,
6952 stoploss_mode: StoplossMode::FromSignal,
6953 rules: vec![],
6954 group_override: Some("alpha".into()),
6955 let_remainder_run: false,
6956 entry_geometry: EntryGeometryPolicy::Strict,
6957 };
6958
6959 let raw_signals = vec![
6960 RawSignal::Entry {
6961 ts: ts(10, 0, 0),
6962 symbol: "EURUSD".into(),
6963 side: Side::Buy,
6964 order_type: OrderType::Market,
6965 price: Some(1.0850),
6966 risk_multiplier: 1.0,
6967 stoploss: None,
6968 targets: vec![1.0910],
6969 group: None,
6970 trade_id: Some("t1".into()),
6971 entry_class: None,
6972 },
6973 RawSignal::Entry {
6974 ts: ts(10, 0, 1),
6975 symbol: "EURUSD".into(),
6976 side: Side::Buy,
6977 order_type: OrderType::Market,
6978 price: Some(1.0857),
6979 risk_multiplier: 1.0,
6980 stoploss: None,
6981 targets: vec![1.0910],
6982 group: None,
6983 trade_id: Some("t2".into()),
6984 entry_class: None,
6985 },
6986 RawSignal::Close {
6988 ts: ts(10, 0, 2),
6989 position: PositionRef::ByTradeId {
6990 trade_id: "t1".into(),
6991 },
6992 },
6993 ];
6994
6995 let config = BacktestConfig {
6996 initial_balance: 10_000.0,
6997 close_on_finish: true,
6998 ..fixed_lot_config()
6999 };
7000 let runner = BacktestRunner::new(config);
7001 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
7002
7003 assert_eq!(result.total_trades, 2);
7005 for trade in &result.trade_log {
7007 assert_eq!(trade.group.as_deref(), Some("alpha"));
7008 }
7009 }
7010
7011 #[test]
7012 fn merged_feed_manual_close_uses_correct_symbol_quote() {
7013 use crate::data_feed::MarketEvent;
7018 let events = vec![
7019 MarketEvent::Tick {
7020 symbol: "XAUUSD".into(),
7021 ts: ts(10, 0, 0),
7022 bid: 5000.0,
7023 ask: 5001.0,
7024 },
7025 MarketEvent::Tick {
7026 symbol: "GBPJPY".into(),
7027 ts: ts(10, 0, 1),
7028 bid: 210.0,
7029 ask: 211.0,
7030 },
7031 MarketEvent::Tick {
7032 symbol: "XAUUSD".into(),
7033 ts: ts(10, 0, 2),
7034 bid: 5050.0,
7035 ask: 5051.0,
7036 },
7037 MarketEvent::Tick {
7039 symbol: "GBPJPY".into(),
7040 ts: ts(10, 0, 3),
7041 bid: 212.0,
7042 ask: 213.0,
7043 },
7044 ];
7045 let mut feed = VecFeed::new(events);
7046
7047 let raw_signals = vec![
7048 RawSignal::Entry {
7049 ts: ts(10, 0, 0),
7050 symbol: "XAUUSD".into(),
7051 side: Side::Buy,
7052 order_type: OrderType::Market,
7053 price: Some(5000.0),
7054 risk_multiplier: 1.0,
7055 stoploss: None,
7056 targets: vec![],
7057 group: None,
7058 trade_id: Some("xau-1".into()),
7059 entry_class: None,
7060 },
7061 RawSignal::Close {
7063 ts: ts(10, 0, 3),
7064 position: PositionRef::ByTradeId {
7065 trade_id: "xau-1".into(),
7066 },
7067 },
7068 ];
7069
7070 let config = BacktestConfig {
7071 initial_balance: 10_000.0,
7072 close_on_finish: false,
7073 ..fixed_lot_config()
7074 };
7075 let runner = BacktestRunner::new(config);
7076 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
7077
7078 assert_eq!(result.total_trades, 1);
7079 let trade = &result.trade_log[0];
7080 assert_eq!(trade.symbol, "XAUUSD");
7081 assert!(
7083 trade.exit_price > 4000.0,
7084 "Exit price should be XAUUSD (~5050), got {}",
7085 trade.exit_price
7086 );
7087 }
7088
7089 fn long_tick_feed(count: usize) -> VecFeed {
7090 let start = ts(10, 0, 0);
7091 VecFeed::new(
7092 (0..count)
7093 .map(|index| {
7094 tick(
7095 "EURUSD",
7096 1.0848,
7097 1.0850,
7098 start + Duration::milliseconds(index as i64),
7099 )
7100 })
7101 .collect(),
7102 )
7103 }
7104
7105 #[test]
7106 fn legacy_replay_can_be_cancelled_during_event_processing() {
7107 let cancelled = std::cell::Cell::new(false);
7108 let mut feed = long_tick_feed(1_000);
7109 let outcome = BacktestRunner::with_defaults().run_raw_signals_controlled(
7110 &mut feed,
7111 Vec::new(),
7112 None,
7113 || cancelled.get(),
7114 |progress| {
7115 if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
7116 cancelled.set(true);
7117 }
7118 },
7119 );
7120
7121 assert_eq!(outcome.unwrap_err(), ReplayCancelled);
7122 assert!(
7123 feed.remaining() > 0,
7124 "cancellation must stop further replay"
7125 );
7126 }
7127
7128 #[test]
7129 fn future_quote_replay_can_be_cancelled_during_event_processing() {
7130 let cancelled = std::cell::Cell::new(false);
7131 let mut feed = long_tick_feed(1_000);
7132 let runner = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default());
7133 let pending = RawSignal::Entry {
7134 ts: ts(10, 0, 0),
7135 symbol: "EURUSD".into(),
7136 side: Side::Buy,
7137 order_type: OrderType::Limit,
7138 price: Some(1.0),
7139 risk_multiplier: 1.0,
7140 stoploss: None,
7141 targets: Vec::new(),
7142 group: None,
7143 trade_id: Some("cancellation-blocker".into()),
7144 entry_class: None,
7145 };
7146 let outcome = runner.run_raw_signals_controlled(
7147 &mut feed,
7148 vec![pending],
7149 None,
7150 || cancelled.get(),
7151 |progress| {
7152 if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
7153 cancelled.set(true);
7154 }
7155 },
7156 );
7157
7158 assert_eq!(outcome.unwrap_err(), ReplayCancelled);
7159 }
7160
7161 #[test]
7162 fn controlled_replay_progress_is_monotonic_and_reaches_event_total() {
7163 let mut feed = long_tick_feed(600);
7164 let mut updates = Vec::new();
7165 BacktestRunner::with_defaults()
7166 .run_raw_signals_controlled(
7167 &mut feed,
7168 Vec::new(),
7169 None,
7170 || false,
7171 |progress| updates.push(progress),
7172 )
7173 .unwrap();
7174
7175 assert!(updates.len() >= 3);
7176 assert!(updates.windows(2).all(|pair| {
7177 pair[0].processed_events <= pair[1].processed_events
7178 && pair[0].processed_signals <= pair[1].processed_signals
7179 && pair[0].total_events <= pair[1].total_events
7180 && pair[0].total_signals <= pair[1].total_signals
7181 }));
7182 assert_eq!(updates.last().unwrap().processed_events, 600);
7183 assert_eq!(updates.last().unwrap().total_events, 600);
7184 }
7185
7186 #[test]
7187 fn legacy_replay_skips_invalid_crossed_and_reversed_quotes_without_nonfinite_pnl() {
7188 let events = vec![
7189 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7190 tick("EURUSD", f64::NAN, 101.0, ts(10, 0, 1)),
7191 tick("EURUSD", 102.0, 101.0, ts(10, 0, 2)),
7192 tick("EURUSD", 90.0, 90.0, ts(9, 59, 59)),
7193 tick("EURUSD", 110.0, 110.0, ts(10, 0, 3)),
7194 ];
7195 let mut feed = VecFeed::new(events);
7196 let signals = vec![RawSignal::Entry {
7197 ts: ts(10, 0, 0),
7198 symbol: "EURUSD".into(),
7199 side: Side::Buy,
7200 order_type: OrderType::Market,
7201 price: Some(100.0),
7202 risk_multiplier: 1.0,
7203 stoploss: None,
7204 targets: vec![],
7205 group: None,
7206 trade_id: Some("safe-feed".into()),
7207 entry_class: None,
7208 }];
7209
7210 let result =
7211 BacktestRunner::new(fixed_lot_config()).run_raw_signals(&mut feed, signals, None);
7212 assert_eq!(result.trade_log.len(), 1);
7213 assert_eq!(result.trade_log[0].exit_price, 110.0);
7214 assert_eq!(result.trade_log[0].pnl, 10.0);
7215 assert!(result.total_pnl.is_finite());
7216 assert!(result.final_balance.is_finite());
7217 }
7218
7219 #[test]
7220 fn legacy_and_future_profile_replay_share_empty_ratio_target_resolution() {
7221 let profile = ManagementProfile {
7222 name: "equal-target".into(),
7223 target_selection: None,
7224 use_targets: vec![1],
7225 close_ratios: vec![],
7226 target_source: TargetSource::FromSignal,
7227 stoploss_mode: StoplossMode::FromSignal,
7228 rules: vec![],
7229 group_override: None,
7230 let_remainder_run: false,
7231 entry_geometry: EntryGeometryPolicy::Strict,
7232 };
7233 let signals = vec![RawSignal::Entry {
7234 ts: ts(10, 0, 0),
7235 symbol: "EURUSD".into(),
7236 side: Side::Buy,
7237 order_type: OrderType::Market,
7238 price: Some(100.0),
7239 risk_multiplier: 1.0,
7240 stoploss: None,
7241 targets: vec![101.0],
7242 group: None,
7243 trade_id: Some("profile-parity".into()),
7244 entry_class: None,
7245 }];
7246 let events = vec![
7247 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7248 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7249 ];
7250
7251 let mut legacy_feed = VecFeed::new(events.clone());
7252 let legacy = BacktestRunner::new(BacktestConfig {
7253 close_on_finish: false,
7254 ..fixed_lot_config()
7255 })
7256 .run_raw_signals(&mut legacy_feed, signals.clone(), Some(&profile));
7257 let mut future_feed = VecFeed::new(events);
7258 let future = BacktestRunner::new_future(
7259 BacktestConfig {
7260 close_on_finish: false,
7261 ..fixed_lot_config()
7262 },
7263 FutureQuoteConfig::default(),
7264 )
7265 .run_raw_signals_future(&mut future_feed, signals, Some(&profile));
7266
7267 assert_eq!(legacy.trade_log.len(), 1);
7268 assert_eq!(future.trade_log.len(), 1);
7269 assert_eq!(legacy.trade_log[0].close_reason, CloseReason::Target);
7270 assert_eq!(future.trade_log[0].close_reason, CloseReason::Target);
7271 assert_eq!(legacy.trade_log[0].size, future.trade_log[0].size);
7272 }
7273
7274 #[test]
7275 fn future_batch_sizes_from_shared_conversion_before_primary_and_uses_primary_eod() {
7276 let currency_plan = RunCurrencyPlan::new(
7277 "USD",
7278 ["EURUSD".to_owned()].into_iter().collect(),
7279 ["EURUSD".to_owned()].into_iter().collect(),
7280 [("EURUSD".to_owned(), "EUR".to_owned())]
7281 .into_iter()
7282 .collect(),
7283 [(
7284 "EUR".to_owned(),
7285 ConversionRoute::Direct {
7286 pair: FxPair {
7287 symbol: "EURUSD".to_owned(),
7288 base_currency: "EUR".to_owned(),
7289 quote_currency: "USD".to_owned(),
7290 },
7291 },
7292 )]
7293 .into_iter()
7294 .collect(),
7295 Vec::new(),
7296 )
7297 .unwrap();
7298 let mut config = fixed_lot_config();
7299 config.sizing = Some(SizingPolicy::FixedRiskAmount { amount: 12.0 });
7300 let future = FutureQuoteConfig {
7301 currency_plan: Some(currency_plan),
7302 conversion_stale_after_ms: 1_000,
7303 ..FutureQuoteConfig::default()
7304 };
7305 let events = vec![
7306 FeedEvent::new(
7307 tick("EURUSD", 1.1, 1.2, ts(10, 0, 0)),
7308 EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
7309 ),
7310 FeedEvent::new(
7311 tick("EURUSD", 2.0, 2.1, ts(10, 0, 1)),
7312 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
7313 ),
7314 ];
7315 let signals = vec![RawSignal::Entry {
7316 ts: ts(10, 0, 0),
7317 symbol: "EURUSD".into(),
7318 side: Side::Buy,
7319 order_type: OrderType::Market,
7320 price: Some(1.0),
7321 risk_multiplier: 1.0,
7322 stoploss: Some(1.19),
7323 targets: Vec::new(),
7324 group: None,
7325 trade_id: Some("shared-conversion".into()),
7326 entry_class: None,
7327 }];
7328
7329 let mut feed = VecFeed::from_feed_events(events);
7330 let result = BacktestRunner::new_future(config, future)
7331 .run_raw_signals_future(&mut feed, signals, None);
7332
7333 assert_eq!(result.recorded_fills.len(), 2);
7334 assert!((result.recorded_fills[0].fill.price - 1.2).abs() < 1.0e-12);
7335 assert!((result.recorded_fills[0].size - 10.0).abs() < 1.0e-12);
7336 assert_eq!(result.recorded_fills[1].execution_ts, Some(ts(10, 0, 0)));
7337 assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 0));
7338 assert!(
7339 result
7340 .mtm_equity_curve
7341 .iter()
7342 .all(|point| point.ts == ts(10, 0, 0))
7343 );
7344 }
7345
7346 #[test]
7347 fn conversion_only_batch_revalues_but_defers_execution_to_primary_quote() {
7348 let currency_plan = RunCurrencyPlan::new(
7349 "USD",
7350 ["EURUSD".to_owned()].into_iter().collect(),
7351 ["EURUSD".to_owned()].into_iter().collect(),
7352 [("EURUSD".to_owned(), "EUR".to_owned())]
7353 .into_iter()
7354 .collect(),
7355 [(
7356 "EUR".to_owned(),
7357 ConversionRoute::Direct {
7358 pair: FxPair {
7359 symbol: "EURUSD".to_owned(),
7360 base_currency: "EUR".to_owned(),
7361 quote_currency: "USD".to_owned(),
7362 },
7363 },
7364 )]
7365 .into_iter()
7366 .collect(),
7367 Vec::new(),
7368 )
7369 .unwrap();
7370 let config = BacktestConfig {
7371 close_on_finish: false,
7372 ..fixed_lot_config()
7373 };
7374 let future = FutureQuoteConfig {
7375 currency_plan: Some(currency_plan),
7376 conversion_stale_after_ms: 10_000,
7377 ..FutureQuoteConfig::default()
7378 };
7379 let events = vec![
7380 FeedEvent::new(
7381 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7382 EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
7383 ),
7384 FeedEvent::new(
7385 tick("EURUSD", 2.0, 2.0, ts(10, 0, 1)),
7386 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
7387 ),
7388 FeedEvent::new(
7389 tick("EURUSD", 110.0, 110.0, ts(10, 0, 2)),
7390 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
7391 ),
7392 ];
7393 let signals = vec![
7394 RawSignal::Entry {
7395 ts: ts(10, 0, 0),
7396 symbol: "EURUSD".into(),
7397 side: Side::Buy,
7398 order_type: OrderType::Market,
7399 price: None,
7400 risk_multiplier: 1.0,
7401 stoploss: None,
7402 targets: Vec::new(),
7403 group: None,
7404 trade_id: Some("conversion-only".into()),
7405 entry_class: None,
7406 },
7407 RawSignal::Close {
7408 ts: ts(10, 0, 1),
7409 position: PositionRef::ByTradeId {
7410 trade_id: "conversion-only".into(),
7411 },
7412 },
7413 ];
7414
7415 let mut feed = VecFeed::from_feed_events(events);
7416 let result = BacktestRunner::new_future(config, future)
7417 .run_raw_signals_future(&mut feed, signals, None);
7418
7419 assert_eq!(result.recorded_fills.len(), 2);
7420 assert_eq!(result.recorded_fills[0].quote_ts, ts(10, 0, 0));
7421 assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 2));
7422 assert!(
7423 result
7424 .mtm_equity_curve
7425 .iter()
7426 .any(|point| point.ts == ts(10, 0, 1))
7427 );
7428 assert_eq!(result.total_pnl, 20.0);
7429 assert_eq!(result.close_events[0].native_pnl, Some(10.0));
7430 assert_eq!(
7431 result.close_events[0]
7432 .pnl_conversion
7433 .as_ref()
7434 .unwrap()
7435 .operation_ts,
7436 ts(10, 0, 2)
7437 );
7438 }
7439
7440 #[test]
7441 fn exact_timestamp_close_updates_balance_before_later_risk_entry() {
7442 let mut config = fixed_lot_config();
7443 config.close_on_finish = false;
7444 config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
7445 let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
7446 spec.digits = 2;
7447 spec.pip_position = 2;
7448 spec.lot_base_units = 1;
7449 spec.lot_step_units = 1;
7450 let future = FutureQuoteConfig {
7451 currency_plan: Some(identity_currency_plan("EURUSD")),
7452 ..FutureQuoteConfig::default()
7453 };
7454 let signals = vec![
7455 RawSignal::Entry {
7456 ts: ts(10, 0, 0),
7457 symbol: "EURUSD".into(),
7458 side: Side::Buy,
7459 order_type: OrderType::Market,
7460 price: None,
7461 risk_multiplier: 1.0,
7462 stoploss: Some(99.0),
7463 targets: Vec::new(),
7464 group: None,
7465 trade_id: Some("first".into()),
7466 entry_class: None,
7467 },
7468 RawSignal::Close {
7469 ts: ts(10, 0, 1),
7470 position: PositionRef::ByTradeId {
7471 trade_id: "first".into(),
7472 },
7473 },
7474 RawSignal::Entry {
7475 ts: ts(10, 0, 1),
7476 symbol: "EURUSD".into(),
7477 side: Side::Buy,
7478 order_type: OrderType::Market,
7479 price: None,
7480 risk_multiplier: 1.0,
7481 stoploss: Some(100.0),
7482 targets: Vec::new(),
7483 group: None,
7484 trade_id: Some("second".into()),
7485 entry_class: None,
7486 },
7487 ];
7488 let mut feed = VecFeed::new(vec![
7489 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7490 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7491 ]);
7492
7493 let result = BacktestRunner::new_future(config, future)
7494 .run_raw_signals_future(&mut feed, signals, None);
7495
7496 assert!((result.total_pnl - 100.0).abs() < 1.0e-12);
7497 assert_eq!(result.open_position_snapshots.len(), 1);
7498 assert_eq!(
7499 result.open_position_snapshots[0].trade_id.as_deref(),
7500 Some("second")
7501 );
7502 assert!((result.open_position_snapshots[0].remaining_size - 101.0).abs() < 1.0e-12);
7503 }
7504
7505 #[test]
7506 fn pending_fill_keeps_placement_size_after_balance_changes() {
7507 let mut config = fixed_lot_config();
7508 config.close_on_finish = false;
7509 config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
7510 let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
7511 spec.digits = 2;
7512 spec.pip_position = 2;
7513 spec.lot_base_units = 1;
7514 spec.lot_step_units = 1;
7515 let future = FutureQuoteConfig {
7516 currency_plan: Some(identity_currency_plan("EURUSD")),
7517 market_entry_sizing_basis: MarketEntrySizingBasis::SignalEntryPrice,
7518 ..FutureQuoteConfig::default()
7519 };
7520 let signals = vec![
7521 RawSignal::Entry {
7522 ts: ts(10, 0, 0),
7523 symbol: "EURUSD".into(),
7524 side: Side::Buy,
7525 order_type: OrderType::Market,
7526 price: None,
7527 risk_multiplier: 1.0,
7528 stoploss: Some(99.0),
7529 targets: Vec::new(),
7530 group: None,
7531 trade_id: Some("market".into()),
7532 entry_class: None,
7533 },
7534 RawSignal::Entry {
7535 ts: ts(10, 0, 0),
7536 symbol: "EURUSD".into(),
7537 side: Side::Buy,
7538 order_type: OrderType::Limit,
7539 price: Some(99.0),
7540 risk_multiplier: 1.0,
7541 stoploss: Some(98.0),
7542 targets: Vec::new(),
7543 group: None,
7544 trade_id: Some("pending".into()),
7545 entry_class: None,
7546 },
7547 RawSignal::Close {
7548 ts: ts(10, 0, 1),
7549 position: PositionRef::ByTradeId {
7550 trade_id: "market".into(),
7551 },
7552 },
7553 ];
7554 let mut feed = VecFeed::new(vec![
7555 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7556 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7557 tick("EURUSD", 99.0, 99.0, ts(10, 0, 2)),
7558 ]);
7559
7560 let result = BacktestRunner::new_future(config, future)
7561 .run_raw_signals_future(&mut feed, signals, None);
7562
7563 assert_eq!(result.pending_order_snapshots.len(), 0);
7564 assert_eq!(result.open_position_snapshots.len(), 1);
7565 assert_eq!(
7566 result.open_position_snapshots[0].trade_id.as_deref(),
7567 Some("pending")
7568 );
7569 assert!((result.open_position_snapshots[0].remaining_size - 100.0).abs() < 1.0e-12);
7570 let metadata = result.execution_metadata.as_ref().unwrap();
7571 assert_eq!(metadata.market_entry_sizing.len(), 1);
7572 assert_eq!(
7573 metadata.market_entry_sizing[0].trade_id.as_deref(),
7574 Some("market")
7575 );
7576 assert_eq!(metadata.entry_profile_resolutions.len(), 2);
7577 assert!(metadata.entry_profile_resolutions.iter().any(|audit| {
7578 audit.trade_id.as_deref() == Some("pending")
7579 && audit.resolution_stage == EntryResolutionStage::PendingPlacement
7580 }));
7581 }
7582
7583 #[test]
7584 fn raw_entries_require_sizing_but_management_only_replay_does_not() {
7585 let entry = RawSignal::Entry {
7586 ts: ts(10, 0, 0),
7587 symbol: "EURUSD".into(),
7588 side: Side::Buy,
7589 order_type: OrderType::Market,
7590 price: None,
7591 risk_multiplier: 1.0,
7592 stoploss: None,
7593 targets: Vec::new(),
7594 group: None,
7595 trade_id: None,
7596 entry_class: None,
7597 };
7598 let mut entry_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
7599 let rejected =
7600 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
7601 .run_raw_signals_future(&mut entry_feed, vec![entry], None);
7602 assert!(rejected.action_dispositions.iter().any(|disposition| {
7603 disposition.action_id == "configuration"
7604 && disposition
7605 .reason
7606 .as_deref()
7607 .is_some_and(|reason| reason.contains("BacktestConfig.sizing"))
7608 }));
7609
7610 let mut management_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
7611 let management =
7612 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
7613 .run_raw_signals_future(
7614 &mut management_feed,
7615 vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
7616 None,
7617 );
7618 assert!(
7619 management
7620 .action_dispositions
7621 .iter()
7622 .all(|disposition| disposition.action_id != "configuration")
7623 );
7624 }
7625
7626 #[test]
7627 fn server_filter_signals_before_market_window() {
7628 let events = vec![
7634 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
7635 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
7636 ];
7637 let mut feed = VecFeed::new(events);
7638
7639 let raw_signals = vec![RawSignal::Entry {
7641 ts: NaiveDate::from_ymd_opt(2026, 1, 1)
7642 .unwrap()
7643 .and_hms_opt(0, 0, 0)
7644 .unwrap(),
7645 symbol: "EURUSD".into(),
7646 side: Side::Buy,
7647 order_type: OrderType::Market,
7648 price: Some(1.0850),
7649 risk_multiplier: 1.0,
7650 stoploss: None,
7651 targets: vec![1.0900],
7652 group: None,
7653 trade_id: None,
7654 entry_class: None,
7655 }];
7656
7657 let runner = BacktestRunner::new(fixed_lot_config());
7658 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
7659
7660 assert_eq!(result.total_trades, 1);
7662 }
7663
7664 #[test]
7665 fn future_replay_audits_profile_resolution_rejection() {
7666 let profile = ManagementProfile {
7667 name: "requires_stop".into(),
7668 target_selection: Some(crate::profile::TargetSelection::None),
7669 use_targets: vec![],
7670 close_ratios: vec![],
7671 target_source: TargetSource::FromSignal,
7672 stoploss_mode: StoplossMode::FromSignalDistance { multiplier: 1.5 },
7673 rules: vec![],
7674 group_override: None,
7675 let_remainder_run: true,
7676 entry_geometry: EntryGeometryPolicy::Strict,
7677 };
7678 let profiles = PreparedEntryProfiles::try_new(
7679 Some(profile),
7680 Vec::<(String, ManagementProfile)>::new(),
7681 )
7682 .unwrap();
7683 let signals = vec![RawSignal::Entry {
7684 ts: ts(10, 0, 0),
7685 symbol: "EURUSD".into(),
7686 side: Side::Buy,
7687 order_type: OrderType::Market,
7688 price: Some(1.1000),
7689 risk_multiplier: 1.0,
7690 stoploss: None,
7691 targets: vec![],
7692 group: None,
7693 trade_id: Some("missing-stop".into()),
7694 entry_class: None,
7695 }];
7696 let mut feed = VecFeed::new(vec![tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1))]);
7697 let result = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default())
7698 .with_entry_profiles(profiles)
7699 .run_raw_signals_future(&mut feed, signals, None);
7700 let audit = &result
7701 .execution_metadata
7702 .as_ref()
7703 .unwrap()
7704 .entry_profile_resolutions[0];
7705 assert_eq!(audit.outcome, ActionDispositionStatus::Rejected);
7706 assert_eq!(audit.rejection_stage.as_deref(), Some("profile_resolution"));
7707 assert!(audit.reason.as_deref().unwrap().contains("signal stoploss"));
7708 }
7709
7710 #[test]
7711 fn future_replay_routes_entry_profiles_and_audits_resolved_levels() {
7712 let default_profile = ManagementProfile {
7713 name: "default".into(),
7714 target_selection: Some(crate::profile::TargetSelection::None),
7715 use_targets: vec![],
7716 close_ratios: vec![],
7717 target_source: TargetSource::FromSignal,
7718 stoploss_mode: StoplossMode::FromSignal,
7719 rules: vec![],
7720 group_override: None,
7721 let_remainder_run: true,
7722 entry_geometry: EntryGeometryPolicy::Strict,
7723 };
7724 let expanded_profile = ManagementProfile {
7725 name: "expanded".into(),
7726 target_selection: None,
7727 use_targets: vec![],
7728 close_ratios: vec![1.0],
7729 target_source: TargetSource::StopDistanceMultiples {
7730 multiples: vec![1.0],
7731 },
7732 stoploss_mode: StoplossMode::FromSignalDistance { multiplier: 1.5 },
7733 rules: vec![],
7734 group_override: None,
7735 let_remainder_run: false,
7736 entry_geometry: EntryGeometryPolicy::Strict,
7737 };
7738 let profiles = PreparedEntryProfiles::try_new(
7739 Some(default_profile),
7740 [("expanded".to_owned(), expanded_profile)],
7741 )
7742 .unwrap();
7743 let signals = vec![
7744 RawSignal::Entry {
7745 ts: ts(10, 0, 0),
7746 symbol: "EURUSD".into(),
7747 side: Side::Buy,
7748 order_type: OrderType::Market,
7749 price: Some(1.1000),
7750 risk_multiplier: 1.0,
7751 stoploss: Some(1.0990),
7752 targets: vec![],
7753 group: Some("same-group".into()),
7754 trade_id: Some("default-entry".into()),
7755 entry_class: None,
7756 },
7757 RawSignal::Entry {
7758 ts: ts(10, 0, 1),
7759 symbol: "EURUSD".into(),
7760 side: Side::Buy,
7761 order_type: OrderType::Market,
7762 price: Some(1.1000),
7763 risk_multiplier: 1.0,
7764 stoploss: Some(1.0990),
7765 targets: vec![],
7766 group: Some("same-group".into()),
7767 trade_id: Some("expanded-entry".into()),
7768 entry_class: Some("expanded".into()),
7769 },
7770 ];
7771 let mut feed = VecFeed::new(vec![
7772 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
7773 tick("EURUSD", 1.1002, 1.1002, ts(10, 0, 2)),
7774 tick("EURUSD", 1.1003, 1.1003, ts(10, 0, 3)),
7775 ]);
7776 let result = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default())
7777 .with_entry_profiles(profiles)
7778 .run_raw_signals_future(&mut feed, signals, None);
7779
7780 let audits = &result
7781 .execution_metadata
7782 .as_ref()
7783 .unwrap()
7784 .entry_profile_resolutions;
7785 assert_eq!(audits.len(), 2);
7786 assert_eq!(
7787 audits[0].selection_source,
7788 EntryProfileSelectionSource::RunDefault
7789 );
7790 assert_eq!(audits[0].selected_profile_name.as_deref(), Some("default"));
7791 assert_eq!(
7792 audits[1].selection_source,
7793 EntryProfileSelectionSource::Mapped
7794 );
7795 assert_eq!(audits[1].entry_class.as_deref(), Some("expanded"));
7796 assert_eq!(audits[1].selected_profile_name.as_deref(), Some("expanded"));
7797 let levels = audits[1].level_resolution.as_ref().unwrap();
7798 let reference = audits[1].level_reference_price.unwrap();
7799 let stop = levels.resolved_stoploss.unwrap();
7800 let target = levels.resolved_targets[0];
7801 let risk_distance = reference - stop;
7802 assert!(risk_distance > 0.0);
7803 assert!((target - reference - risk_distance).abs() < 1.0e-9);
7804 }
7805}