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