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