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