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, 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::Strategy;
49
50#[derive(Debug, Clone)]
53pub struct FutureQuoteConfig {
54 pub signal_latency_ms: i64,
56 pub slippage_pips: f64,
58 pub stale_quote_after_ms: Option<i64>,
60 pub pnl_epsilon: f64,
62 pub currency_plan: Option<RunCurrencyPlan>,
64 pub conversion_stale_after_ms: i64,
66 pub mtm_output: MtmOutputPolicy,
68}
69
70impl Default for FutureQuoteConfig {
71 fn default() -> Self {
72 Self {
73 signal_latency_ms: 0,
74 slippage_pips: 0.0,
75 stale_quote_after_ms: None,
76 pnl_epsilon: 1.0e-9,
77 currency_plan: None,
78 conversion_stale_after_ms: 300_000,
79 mtm_output: MtmOutputPolicy::default(),
80 }
81 }
82}
83
84#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
86pub struct ReplayProgress {
87 pub processed_events: usize,
88 pub total_events: usize,
89 pub processed_signals: usize,
90 pub total_signals: usize,
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
95#[error("backtest replay cancelled")]
96pub struct ReplayCancelled;
97
98#[derive(Debug, thiserror::Error)]
100pub enum StreamingReplayError<E> {
101 #[error("market-data stream failed: {0}")]
102 Feed(E),
103 #[error(transparent)]
104 Cancelled(#[from] ReplayCancelled),
105}
106
107const REPLAY_PROGRESS_INTERVAL: usize = 256;
108
109fn should_report_progress(processed: usize, total: usize) -> bool {
110 processed == total || processed.is_multiple_of(REPLAY_PROGRESS_INTERVAL)
111}
112
113#[derive(Debug, Clone)]
114struct ScheduledSignal {
115 sequence: u64,
116 signal_ts: NaiveDateTime,
117 effective_ts: NaiveDateTime,
118 signal: RawSignal,
119}
120
121#[derive(Debug, Clone)]
122struct QueuedAction {
123 action_id: String,
124 action_kind: String,
125 action: Action,
126 execution: Option<ExecutionFill>,
127 symbol: String,
128 signal_ts: NaiveDateTime,
129 effective_ts: NaiveDateTime,
130 entry_signal: Option<RawSignal>,
131 entry_profile: Option<ManagementProfile>,
132}
133
134#[derive(Debug, thiserror::Error)]
135enum FutureTransactionError {
136 #[error(transparent)]
137 Core(#[from] FutureApplyError),
138 #[error(transparent)]
139 Accounting(#[from] FutureExecutorError),
140}
141
142#[derive(Debug, Clone, Copy, PartialEq, Eq)]
144pub enum EquityObservationKind {
145 PreSettlement,
146 PostOutput,
147 ConversionRevaluation,
148 QuiescentTermination,
149 EndOfData,
150}
151
152impl EquityObservationKind {
153 pub const fn as_str(self) -> &'static str {
154 match self {
155 Self::PreSettlement => "pre_settlement",
156 Self::PostOutput => "post_output",
157 Self::ConversionRevaluation => "conversion_revaluation",
158 Self::QuiescentTermination => "quiescent_termination",
159 Self::EndOfData => "end_of_data",
160 }
161 }
162}
163
164struct BufferedFutureFeed {
165 events: VecDeque<FeedEvent>,
166 total_events: usize,
167}
168
169impl BufferedFutureFeed {
170 fn new(events: Vec<FeedEvent>) -> Self {
171 let total_events = events.len();
172 Self {
173 events: VecDeque::from(events),
174 total_events,
175 }
176 }
177}
178
179impl DataFeed for BufferedFutureFeed {
180 fn next_event(&mut self) -> Option<MarketEvent> {
181 self.events.pop_front().map(|event| event.event)
182 }
183
184 fn peek(&self) -> Option<&MarketEvent> {
185 self.events.front().map(|event| &event.event)
186 }
187
188 fn next_batch(&mut self) -> Option<TimestampBatch> {
189 let ts = self.events.front()?.event.ts();
190 let mut events = Vec::new();
191 while self
192 .events
193 .front()
194 .is_some_and(|event| event.event.ts() == ts)
195 {
196 events.push(self.events.pop_front().expect("front checked"));
197 }
198 Some(TimestampBatch { ts, events })
199 }
200
201 fn total_events(&self) -> Option<usize> {
202 Some(self.total_events)
203 }
204}
205
206struct DataFeedBatchAdapter<'a, F> {
207 feed: &'a mut F,
208}
209
210impl<F: DataFeed> FallibleBatchFeed for DataFeedBatchAdapter<'_, F> {
211 type Error = Infallible;
212
213 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
214 Ok(self.feed.next_batch())
215 }
216}
217
218#[derive(Debug, Clone)]
220pub struct BacktestConfig {
221 pub initial_balance: f64,
223 pub close_on_finish: bool,
226 pub fill_model: FillModel,
231 pub contract_sizes: HashMap<String, f64>,
240 pub sizing: Option<SizingPolicy>,
243 pub symbol_specs: HashMap<String, qs_symbols::SymbolSpec>,
246 pub instrument_manifest: Option<ReplayInstrumentManifest>,
248}
249
250impl Default for BacktestConfig {
251 fn default() -> Self {
252 Self {
253 initial_balance: 10_000.0,
254 close_on_finish: true,
255 fill_model: FillModel::default(),
256 contract_sizes: HashMap::new(),
257 sizing: None,
258 symbol_specs: HashMap::new(),
259 instrument_manifest: None,
260 }
261 }
262}
263
264pub struct BacktestRunner {
266 engine: TradeEngine,
267 executor: BacktestExecutor,
268 config: BacktestConfig,
269 future_config: Option<FutureQuoteConfig>,
270 evaluation_options: EvaluationOptions,
271 instrument_sizing: Vec<InstrumentSizingArtifact>,
272}
273
274impl BacktestRunner {
275 pub fn new(config: BacktestConfig) -> Self {
277 let executor =
278 BacktestExecutor::new(config.initial_balance, effective_contract_sizes(&config));
279 Self {
280 engine: TradeEngine::with_fill_model(config.fill_model),
281 executor,
282 config,
283 future_config: None,
284 evaluation_options: EvaluationOptions::default(),
285 instrument_sizing: Vec::new(),
286 }
287 }
288
289 pub fn new_future(config: BacktestConfig, future_config: FutureQuoteConfig) -> Self {
291 let executor =
292 BacktestExecutor::new(config.initial_balance, effective_contract_sizes(&config));
293 let engine = TradeEngine::with_fill_model_and_deterministic_ids(config.fill_model);
294 Self {
295 engine,
296 executor,
297 config,
298 future_config: Some(future_config),
299 evaluation_options: EvaluationOptions::default(),
300 instrument_sizing: Vec::new(),
301 }
302 }
303
304 pub fn with_defaults() -> Self {
306 Self::new(BacktestConfig::default())
307 }
308
309 pub fn with_evaluation_options(mut self, options: EvaluationOptions) -> Self {
312 self.evaluation_options = options;
313 self
314 }
315
316 pub fn engine(&self) -> &TradeEngine {
318 &self.engine
319 }
320
321 pub fn executor(&self) -> &BacktestExecutor {
323 &self.executor
324 }
325
326 pub fn run_strategy<F: DataFeed, S: Strategy>(
340 mut self,
341 feed: &mut F,
342 strategy: &mut S,
343 ) -> BacktestResult {
344 if validate_replay_config(&self.config, None, &[]).is_err() {
345 return rejected_legacy_result(&self.config);
346 }
347 let mut last_quote_ts = BTreeMap::new();
348 while let Some(event) = feed.next_event() {
349 let quote = event.to_quote();
350 if !accept_legacy_quote("e, &mut last_quote_ts) {
351 continue;
352 }
353
354 let price_effects = self.engine.on_price("e);
356 self.executor
357 .process_effects(&price_effects, &self.engine, "e);
358
359 let actions = strategy.on_event(&event);
361 self.apply_actions(actions, "e);
362 }
363
364 let final_actions = strategy.on_finished();
366 if !final_actions.is_empty() {
367 if let Some(last_quote) = self.last_available_quote() {
371 self.apply_actions(final_actions, &last_quote);
372 let effects = self.engine.on_price(&last_quote);
374 self.executor
375 .process_effects(&effects, &self.engine, &last_quote);
376 }
377 }
378
379 self.close_remaining_if_configured();
381
382 BacktestResult::from_trade_log(self.config.initial_balance, self.executor.trade_log)
383 }
384
385 fn apply_actions(&mut self, actions: Vec<Action>, quote: &PriceQuote) {
389 for action in actions {
390 self.apply_single_action(action, quote.ts, quote);
391 }
392 }
393
394 fn apply_single_action(
401 &mut self,
402 action: Action,
403 ts: chrono::NaiveDateTime,
404 quote: &PriceQuote,
405 ) {
406 match self.engine.apply_action(action, ts) {
407 Ok(effects) => {
408 self.executor.process_effects(&effects, &self.engine, quote);
409 }
410 Err(_) => {
411 }
415 }
416 }
417
418 fn last_available_quote(&self) -> Option<PriceQuote> {
420 for pos in self.engine.open_positions() {
423 if let Some(q) = self.engine.last_quote(&pos.data.symbol) {
424 return Some(q.clone());
425 }
426 }
427 for pos in self.engine.closed_positions() {
429 if let Some(q) = self.engine.last_quote(&pos.data.symbol) {
430 return Some(q.clone());
431 }
432 }
433 None
434 }
435
436 pub fn run_raw_signals<F: DataFeed>(
445 self,
446 feed: &mut F,
447 raw_signals: Vec<RawSignal>,
448 profile: Option<&ManagementProfile>,
449 ) -> BacktestResult {
450 self.run_raw_signals_controlled(feed, raw_signals, profile, || false, |_| {})
451 .expect("non-cancellable replay cannot be cancelled")
452 }
453
454 pub fn run_raw_signals_controlled<F, C, P>(
460 mut self,
461 feed: &mut F,
462 raw_signals: Vec<RawSignal>,
463 profile: Option<&ManagementProfile>,
464 mut is_cancelled: C,
465 mut on_progress: P,
466 ) -> std::result::Result<BacktestResult, ReplayCancelled>
467 where
468 F: DataFeed,
469 C: FnMut() -> bool,
470 P: FnMut(ReplayProgress),
471 {
472 if let Some(future_config) = self.future_config.clone() {
473 return self.run_raw_signals_future_controlled(
474 feed,
475 raw_signals,
476 profile,
477 future_config,
478 &mut is_cancelled,
479 &mut on_progress,
480 );
481 }
482 if validate_replay_config(&self.config, None, &raw_signals).is_err()
483 || profile.is_some_and(|profile| profile.validate().is_err())
484 {
485 return Ok(rejected_legacy_result(&self.config));
486 }
487
488 let total_events = feed.total_events().unwrap_or(0);
489 let total_signals = raw_signals.len();
490 let mut processed_events = 0;
491 let mut sig_idx = 0;
492 let mut last_quote_ts = BTreeMap::new();
493 on_progress(ReplayProgress {
494 processed_events,
495 total_events,
496 processed_signals: sig_idx,
497 total_signals,
498 });
499
500 while let Some(event) = feed.next_event() {
501 if is_cancelled() {
502 return Err(ReplayCancelled);
503 }
504 let quote = event.to_quote();
505 if !accept_legacy_quote("e, &mut last_quote_ts) {
506 processed_events += 1;
507 if should_report_progress(processed_events, total_events) {
508 on_progress(ReplayProgress {
509 processed_events,
510 total_events,
511 processed_signals: sig_idx,
512 total_signals,
513 });
514 }
515 continue;
516 }
517
518 while sig_idx < raw_signals.len() && raw_signals[sig_idx].ts() <= event.ts() {
520 if is_cancelled() {
521 return Err(ReplayCancelled);
522 }
523 self.process_raw_signal(&raw_signals[sig_idx], profile, "e);
524 sig_idx += 1;
525 if should_report_progress(sig_idx, total_signals) {
526 on_progress(ReplayProgress {
527 processed_events,
528 total_events,
529 processed_signals: sig_idx,
530 total_signals,
531 });
532 }
533 }
534
535 let effects = self.engine.on_price("e);
537 self.executor
538 .process_effects(&effects, &self.engine, "e);
539 processed_events += 1;
540 if should_report_progress(processed_events, total_events) {
541 on_progress(ReplayProgress {
542 processed_events,
543 total_events,
544 processed_signals: sig_idx,
545 total_signals,
546 });
547 }
548 }
549
550 if sig_idx < raw_signals.len()
552 && let Some(last_quote) = self.last_available_quote()
553 {
554 while sig_idx < raw_signals.len() {
555 if is_cancelled() {
556 return Err(ReplayCancelled);
557 }
558 self.process_raw_signal(&raw_signals[sig_idx], profile, &last_quote);
559 sig_idx += 1;
560 if should_report_progress(sig_idx, total_signals) {
561 on_progress(ReplayProgress {
562 processed_events,
563 total_events,
564 processed_signals: sig_idx,
565 total_signals,
566 });
567 }
568 }
569 let effects = self.engine.on_price(&last_quote);
571 self.executor
572 .process_effects(&effects, &self.engine, &last_quote);
573 }
574
575 if is_cancelled() {
576 return Err(ReplayCancelled);
577 }
578
579 self.close_remaining_if_configured();
581 on_progress(ReplayProgress {
582 processed_events,
583 total_events,
584 processed_signals: sig_idx,
585 total_signals,
586 });
587
588 Ok(BacktestResult::from_trade_log(
589 self.config.initial_balance,
590 self.executor.trade_log,
591 ))
592 }
593
594 pub fn run_raw_signals_future<F: DataFeed>(
600 self,
601 feed: &mut F,
602 raw_signals: Vec<RawSignal>,
603 profile: Option<&ManagementProfile>,
604 ) -> BacktestResult {
605 let future_config = self.future_config.clone().unwrap_or_default();
606 self.run_raw_signals_future_with_config(feed, raw_signals, profile, future_config)
607 }
608
609 #[allow(clippy::too_many_arguments)]
613 pub fn run_raw_signals_future_streaming_controlled<F, C, P>(
614 mut self,
615 feed: &mut F,
616 primary_eod: Option<NaiveDateTime>,
617 raw_signals: Vec<RawSignal>,
618 profile: Option<&ManagementProfile>,
619 mut is_cancelled: C,
620 mut on_progress: P,
621 ) -> std::result::Result<BacktestResult, StreamingReplayError<F::Error>>
622 where
623 F: FallibleBatchFeed,
624 C: FnMut() -> bool,
625 P: FnMut(ReplayProgress),
626 {
627 let future = self.future_config.clone().unwrap_or_default();
628 self.future_config = Some(future.clone());
629 if let Err(error) = validate_replay_config(&self.config, Some(&future), &raw_signals) {
630 return Ok(rejected_future_result(
631 &self.config,
632 &future,
633 self.evaluation_options,
634 error,
635 ));
636 }
637 if let Some(profile) = profile
638 && let Err(error) = profile.validate()
639 {
640 return Ok(rejected_future_result(
641 &self.config,
642 &future,
643 self.evaluation_options,
644 error.to_string(),
645 ));
646 }
647
648 self.run_raw_signals_future_batches(
649 feed,
650 primary_eod,
651 raw_signals,
652 profile,
653 future,
654 None,
655 0,
656 0,
657 &mut is_cancelled,
658 &mut on_progress,
659 )
660 }
661
662 pub fn run_raw_signals_future_fallible<F: FallibleBatchFeed>(
666 self,
667 feed: &mut F,
668 raw_signals: Vec<RawSignal>,
669 profile: Option<&ManagementProfile>,
670 ) -> Result<BacktestResult, F::Error> {
671 let future_config = self.future_config.clone().unwrap_or_default();
672 self.run_raw_signals_future_fallible_with_config(feed, raw_signals, profile, future_config)
673 }
674
675 pub fn run_raw_signals_future_fallible_with_config<F: FallibleBatchFeed>(
677 self,
678 feed: &mut F,
679 raw_signals: Vec<RawSignal>,
680 profile: Option<&ManagementProfile>,
681 future: FutureQuoteConfig,
682 ) -> Result<BacktestResult, F::Error> {
683 if let Err(error) = validate_replay_config(&self.config, Some(&future), &raw_signals) {
684 return Ok(rejected_future_result(
685 &self.config,
686 &future,
687 self.evaluation_options,
688 error,
689 ));
690 }
691 if let Some(profile) = profile
692 && let Err(error) = profile.validate()
693 {
694 return Ok(rejected_future_result(
695 &self.config,
696 &future,
697 self.evaluation_options,
698 error.to_string(),
699 ));
700 }
701
702 let mut events = Vec::new();
703 while let Some(batch) = FallibleBatchFeed::next_batch(feed)? {
704 events.extend(batch.events);
705 }
706 let mut buffered_feed = BufferedFutureFeed::new(events);
707 Ok(self.run_raw_signals_future_with_config(
708 &mut buffered_feed,
709 raw_signals,
710 profile,
711 future,
712 ))
713 }
714
715 fn process_raw_signal(
718 &mut self,
719 signal: &RawSignal,
720 profile: Option<&ManagementProfile>,
721 quote: &PriceQuote,
722 ) {
723 let ts = signal.ts();
724
725 if signal.is_entry() {
726 let mut signal = signal.clone();
727 if let RawSignal::Entry {
728 side,
729 order_type: OrderType::Market,
730 price,
731 ..
732 } = &mut signal
733 && price.is_none()
734 {
735 *price = Some(match side {
736 Side::Buy => quote.ask,
737 Side::Sell => quote.bid,
738 });
739 }
740 let resolved = match profile {
741 Some(profile) => profile.apply_entry_signal(&signal),
742 None => resolve_unprofiled_entry(&signal),
743 };
744 if let Ok(Some(resolved)) = resolved
745 && let Ok(action) =
746 self.finalize_resolved_entry(resolved, self.executor.balance, ts, None)
747 {
748 self.apply_single_action(action, ts, quote);
749 }
750 } else {
751 let actions = resolve_signal(signal, &self.engine);
752 for action in actions {
753 self.apply_single_action(action, ts, quote);
754 }
755 }
756 }
757
758 fn finalize_resolved_entry(
759 &mut self,
760 mut resolved: ResolvedEntry,
761 balance_before: f64,
762 operation_ts: NaiveDateTime,
763 conversion_quotes: Option<&ConversionQuoteBook>,
764 ) -> Result<Action, String> {
765 let policy = self
766 .config
767 .sizing
768 .as_ref()
769 .ok_or_else(|| "raw entry requires a sizing policy".to_owned())?;
770 let entry_price = resolved.price.ok_or_else(|| {
771 if resolved.order_type == OrderType::Market {
772 "market entry requires an execution price".to_owned()
773 } else {
774 "pending entry requires a requested price".to_owned()
775 }
776 })?;
777 let explicit_spec = explicit_instrument_spec(&self.config, &resolved.symbol);
778 let legacy_spec = self.config.symbol_specs.get(&resolved.symbol);
779 if explicit_spec.is_none() && legacy_spec.is_none() {
780 return Err(format!(
781 "missing instrument or symbol spec for {}",
782 resolved.symbol
783 ));
784 }
785
786 let (account_loss_per_lot, native_to_account_rate) = if is_monetary_sizing(policy) {
787 let stop = resolved
788 .stoploss
789 .ok_or_else(|| "monetary sizing requires a protective stop".to_owned())?;
790 let native_loss = match explicit_spec {
791 Some(spec) => compute_instrument_native_loss_per_lot(
792 resolved.side,
793 entry_price,
794 stop,
795 u16::from(spec.price.display_scale),
796 &spec.economics,
797 )
798 .map_err(|error| error.to_string())?,
799 None => compute_native_loss_per_lot(
800 resolved.side,
801 entry_price,
802 stop,
803 legacy_spec.expect("legacy spec presence checked"),
804 )
805 .map_err(|error| error.to_string())?,
806 };
807 let plan = self
808 .future_config
809 .as_ref()
810 .and_then(|config| config.currency_plan.as_ref())
811 .ok_or_else(|| "monetary sizing requires a FutureQuote currency plan".to_owned())?;
812 let route = plan
813 .route_for_primary_symbol(&resolved.symbol)
814 .ok_or_else(|| {
815 format!(
816 "currency plan has no frozen route for primary symbol {}",
817 resolved.symbol
818 )
819 })?;
820 let converted = conversion_quotes
821 .ok_or_else(|| "monetary sizing requires conversion quotes".to_owned())?
822 .convert_route(-native_loss, operation_ts, route)
823 .map_err(|error| error.to_string())?;
824 let account_loss = -converted.output_amount;
825 (Some(account_loss), Some(account_loss / native_loss))
826 } else {
827 (None, None)
828 };
829
830 let sizing = match explicit_spec {
831 Some(spec) => compute_instrument_size_for_spec(
832 policy,
833 resolved.risk_multiplier,
834 balance_before,
835 resolved.side,
836 entry_price,
837 resolved.stoploss,
838 spec,
839 native_to_account_rate,
840 )
841 .map_err(|error| error.to_string())?,
842 None => compute_size(
843 policy,
844 resolved.risk_multiplier,
845 balance_before,
846 resolved.side,
847 entry_price,
848 resolved.stoploss,
849 legacy_spec.expect("legacy spec presence checked"),
850 account_loss_per_lot,
851 )
852 .map_err(|error| error.to_string())?,
853 };
854 if let Some(quantity) = sizing.quantity_adjustment {
855 self.instrument_sizing.push(InstrumentSizingArtifact {
856 symbol: resolved.symbol.clone(),
857 operation_ts,
858 quantity,
859 final_notional: sizing.final_notional.clone(),
860 });
861 }
862 let target_steps = allocate_target_steps(
863 sizing.final_lot_steps,
864 &resolved.target_resolution.weights,
865 resolved.target_resolution.remainder,
866 )
867 .map_err(|error| error.to_string())?;
868 if target_steps.len() != resolved.targets.len() {
869 return Err("target allocation does not match resolved targets".to_owned());
870 }
871 for (target, steps) in resolved.targets.iter_mut().zip(target_steps) {
872 target.close_ratio = steps as f64 / sizing.final_lot_steps as f64;
873 }
874
875 Ok(resolved.into_action(sizing.final_lot))
876 }
877
878 fn close_remaining_if_configured(&mut self) {
884 if !self.config.close_on_finish {
885 return;
886 }
887
888 let open_ids: Vec<String> = self
889 .engine
890 .open_positions()
891 .iter()
892 .map(|p| p.data.id.clone())
893 .collect();
894
895 for id in open_ids {
896 let symbol = match self.engine.get_position(&id) {
897 Some(pos) => pos.data.symbol.clone(),
898 None => continue,
899 };
900 let quote = match self.engine.last_quote(&symbol) {
901 Some(q) => q.clone(),
902 None => continue,
903 };
904
905 if let Ok(effects) = self.engine.apply_action(
906 Action::ClosePosition {
907 position_id: id.clone(),
908 },
909 quote.ts,
910 ) {
911 self.executor
912 .process_effects(&effects, &self.engine, "e);
913 }
914 }
915 }
916
917 pub fn run_raw_signals_future_with_config<F: DataFeed>(
924 self,
925 source_feed: &mut F,
926 raw_signals: Vec<RawSignal>,
927 profile: Option<&ManagementProfile>,
928 future: FutureQuoteConfig,
929 ) -> BacktestResult {
930 self.run_raw_signals_future_controlled(
931 source_feed,
932 raw_signals,
933 profile,
934 future,
935 &mut || false,
936 &mut |_| {},
937 )
938 .expect("non-cancellable replay cannot be cancelled")
939 }
940
941 fn run_raw_signals_future_controlled<F, C, P>(
942 mut self,
943 source_feed: &mut F,
944 raw_signals: Vec<RawSignal>,
945 profile: Option<&ManagementProfile>,
946 future: FutureQuoteConfig,
947 is_cancelled: &mut C,
948 on_progress: &mut P,
949 ) -> std::result::Result<BacktestResult, ReplayCancelled>
950 where
951 F: DataFeed,
952 C: FnMut() -> bool,
953 P: FnMut(ReplayProgress),
954 {
955 self.future_config = Some(future.clone());
956 if let Err(error) = validate_replay_config(&self.config, Some(&future), &raw_signals) {
957 return Ok(rejected_future_result(
958 &self.config,
959 &future,
960 self.evaluation_options,
961 error,
962 ));
963 }
964 if let Some(profile) = profile
965 && let Err(error) = profile.validate()
966 {
967 return Ok(rejected_future_result(
968 &self.config,
969 &future,
970 self.evaluation_options,
971 error.to_string(),
972 ));
973 }
974
975 let total_events = source_feed.total_events();
976 let mut processed_events = 0;
977 let mut ordered_events = Vec::<FeedEvent>::new();
978 let mut source_last_ts = BTreeMap::<String, NaiveDateTime>::new();
979 let mut invalid_quotes = 0_u64;
980 while let Some(batch) = source_feed.next_batch() {
981 if is_cancelled() {
982 return Err(ReplayCancelled);
983 }
984 for feed_event in batch.events {
985 if is_cancelled() {
986 return Err(ReplayCancelled);
987 }
988 let symbol = feed_event.event.symbol().to_owned();
989 let event_ts = feed_event.event.ts();
990 if source_last_ts
991 .get(&symbol)
992 .is_some_and(|last| *last > event_ts)
993 {
994 invalid_quotes += 1;
995 processed_events += 1;
996 continue;
997 }
998 source_last_ts.insert(symbol, event_ts);
999 ordered_events.push(feed_event);
1000 }
1001 }
1002 if is_cancelled() {
1003 return Err(ReplayCancelled);
1004 }
1005 ordered_events.sort_by_key(FeedEvent::ordering_key);
1006 if is_cancelled() {
1007 return Err(ReplayCancelled);
1008 }
1009 let primary_eod = ordered_events
1010 .iter()
1011 .filter(|event| event.metadata.roles.primary)
1012 .filter_map(|event| event.event.to_valid_quote())
1013 .map(|quote| quote.ts)
1014 .max();
1015 let mut ordered_feed = crate::data_feed::VecFeed::from_feed_events(ordered_events);
1016 let mut feed = DataFeedBatchAdapter {
1017 feed: &mut ordered_feed,
1018 };
1019 match self.run_raw_signals_future_batches(
1020 &mut feed,
1021 primary_eod,
1022 raw_signals,
1023 profile,
1024 future,
1025 total_events,
1026 processed_events,
1027 invalid_quotes,
1028 is_cancelled,
1029 on_progress,
1030 ) {
1031 Ok(result) => Ok(result),
1032 Err(StreamingReplayError::Cancelled(error)) => Err(error),
1033 Err(StreamingReplayError::Feed(error)) => match error {},
1034 }
1035 }
1036
1037 #[allow(clippy::too_many_arguments)]
1038 fn run_raw_signals_future_batches<F, C, P>(
1039 mut self,
1040 feed: &mut F,
1041 primary_eod: Option<NaiveDateTime>,
1042 raw_signals: Vec<RawSignal>,
1043 profile: Option<&ManagementProfile>,
1044 future: FutureQuoteConfig,
1045 known_total_events: Option<usize>,
1046 mut processed_events: usize,
1047 mut invalid_quotes: u64,
1048 is_cancelled: &mut C,
1049 on_progress: &mut P,
1050 ) -> std::result::Result<BacktestResult, StreamingReplayError<F::Error>>
1051 where
1052 F: FallibleBatchFeed,
1053 C: FnMut() -> bool,
1054 P: FnMut(ReplayProgress),
1055 {
1056 let total_events = known_total_events.unwrap_or(0);
1057 let total_signals = raw_signals.len();
1058 let mut processed_signals = 0;
1059 on_progress(ReplayProgress {
1060 processed_events,
1061 total_events,
1062 processed_signals,
1063 total_signals,
1064 });
1065
1066 let execution_model = ExecutionModel::new(
1067 qs_core::types::ExecutionConvention::FutureQuoteV1,
1068 self.config.fill_model,
1069 if future.slippage_pips == 0.0 {
1070 SlippageModel::None
1071 } else {
1072 SlippageModel::FixedPips {
1073 pips: future.slippage_pips,
1074 }
1075 },
1076 );
1077 let pricer = ExecutionPricer::new(execution_model);
1078 let mut scheduled: Vec<ScheduledSignal> = raw_signals
1079 .into_iter()
1080 .enumerate()
1081 .map(|(sequence, signal)| {
1082 let signal_ts = signal.ts();
1083 ScheduledSignal {
1084 sequence: sequence as u64,
1085 signal_ts,
1086 effective_ts: signal_ts
1087 .checked_add_signed(Duration::milliseconds(future.signal_latency_ms))
1088 .expect("signal latency overflow was validated before scheduling"),
1089 signal,
1090 }
1091 })
1092 .collect();
1093 scheduled.sort_by_key(|signal| (signal.effective_ts, signal.sequence));
1094 let mut scheduled = VecDeque::from(scheduled);
1095 let mut queued = VecDeque::<QueuedAction>::new();
1096 let mut lifecycle = LifecycleLedger::new();
1097 let contract_sizes = effective_contract_sizes(&self.config);
1098 let mut future_executor = FutureExecutor::new(
1099 self.config.initial_balance,
1100 contract_sizes.clone(),
1101 future.pnl_epsilon,
1102 )
1103 .with_currency_plan(future.currency_plan.clone());
1104 let mut portfolio =
1105 PortfolioRecorder::new(self.config.initial_balance, contract_sizes.clone())
1106 .with_fill_model(self.config.fill_model)
1107 .with_stale_quote_after_millis(future.stale_quote_after_ms)
1108 .with_currency_plan(future.currency_plan.clone());
1109 let mut mtm_curve = MtmCurveCollector::new(future.mtm_output)
1110 .expect("MTM output policy was validated before replay");
1111 let mut last_mtm_candidate = None;
1112 let mut conversion_quotes =
1113 ConversionQuoteBook::new(Duration::milliseconds(future.conversion_stale_after_ms))
1114 .expect("conversion quote staleness was validated before replay");
1115 if let Some(plan) = future.currency_plan.as_ref() {
1116 for quote in plan.strict_before_warmup_quotes() {
1117 conversion_quotes
1118 .record_canonical_tick(quote.clone())
1119 .expect("currency plan warmups were validated during construction");
1120 }
1121 }
1122 let mut last_quote_ts = BTreeMap::<String, NaiveDateTime>::new();
1123 let mut last_processed_primary_ts = None;
1124 let mut effective_terminal_ts = primary_eod;
1125 let mut terminated_quiescently = false;
1126
1127 while let Some(mut batch) =
1128 FallibleBatchFeed::next_batch(feed).map_err(StreamingReplayError::Feed)?
1129 {
1130 let batch_ts = batch.ts;
1131 if is_cancelled() {
1132 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1133 }
1134 batch
1135 .events
1136 .sort_by_key(|event| (event.metadata.series_rank, event.metadata.row_sequence));
1137 let mut accepted = Vec::new();
1138 for feed_event in batch.events {
1139 if is_cancelled() {
1140 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1141 }
1142 let quote = feed_event.event.to_quote();
1143 if ExecutionPricer::validate_quote("e).is_err()
1144 || last_quote_ts
1145 .get("e.symbol)
1146 .is_some_and(|last| *last > quote.ts)
1147 {
1148 invalid_quotes += 1;
1149 processed_events += 1;
1150 if should_report_progress(processed_events, total_events) {
1151 on_progress(ReplayProgress {
1152 processed_events,
1153 total_events,
1154 processed_signals,
1155 total_signals,
1156 });
1157 }
1158 continue;
1159 }
1160 last_quote_ts.insert(quote.symbol.clone(), quote.ts);
1161 accepted.push((feed_event, quote));
1162 }
1163
1164 for (feed_event, quote) in &accepted {
1165 if is_cancelled() {
1166 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1167 }
1168 if feed_event.metadata.roles.conversion
1169 && matches!(feed_event.event, MarketEvent::Tick { .. })
1170 && conversion_quotes
1171 .record_canonical_tick(quote.clone())
1172 .is_err()
1173 {
1174 invalid_quotes += 1;
1175 }
1176 }
1177 let valuation_only = accepted
1178 .iter()
1179 .any(|event| event.0.metadata.roles.conversion)
1180 && !accepted.iter().any(|event| event.0.metadata.roles.primary)
1181 && primary_eod.is_some_and(|eod| batch_ts <= eod);
1182
1183 let mut primary_quotes = Vec::new();
1184 let mut batch_quotes = BTreeMap::new();
1185 for (feed_event, quote) in accepted {
1186 if is_cancelled() {
1187 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1188 }
1189 if !feed_event.metadata.roles.primary {
1190 processed_events += 1;
1191 if should_report_progress(processed_events, total_events) {
1192 on_progress(ReplayProgress {
1193 processed_events,
1194 total_events,
1195 processed_signals,
1196 total_signals,
1197 });
1198 }
1199 continue;
1200 }
1201
1202 last_processed_primary_ts = Some(quote.ts);
1203 portfolio.record_quote(quote.clone());
1204 observe_future_equity(
1205 &mut portfolio,
1206 &future_executor,
1207 quote.ts,
1208 &conversion_quotes,
1209 EquityObservationKind::PreSettlement,
1210 &mut mtm_curve,
1211 &mut last_mtm_candidate,
1212 false,
1213 );
1214 batch_quotes.insert(quote.symbol.clone(), quote.clone());
1215 primary_quotes.push(quote);
1216 }
1217
1218 if let Some(representative_quote) = primary_quotes.first() {
1219 let mut settled_quotes = vec![false; primary_quotes.len()];
1220
1221 self.execute_queued_future(
1223 &batch_quotes,
1224 false,
1225 &mut queued,
1226 &mut lifecycle,
1227 &mut future_executor,
1228 &mut portfolio,
1229 &pricer,
1230 &conversion_quotes,
1231 );
1232 let increasing_symbols =
1233 queued_exposure_symbols(&queued, &batch_quotes, representative_quote.ts);
1234 invalid_quotes += self.settle_future_batch_symbols(
1235 &primary_quotes,
1236 &mut settled_quotes,
1237 Some(&increasing_symbols),
1238 &mut lifecycle,
1239 &mut future_executor,
1240 &mut portfolio,
1241 &pricer,
1242 &conversion_quotes,
1243 );
1244 self.execute_queued_future(
1245 &batch_quotes,
1246 true,
1247 &mut queued,
1248 &mut lifecycle,
1249 &mut future_executor,
1250 &mut portfolio,
1251 &pricer,
1252 &conversion_quotes,
1253 );
1254
1255 while scheduled
1256 .front()
1257 .is_some_and(|signal| signal.effective_ts < representative_quote.ts)
1258 {
1259 if is_cancelled() {
1260 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1261 }
1262 let signal = scheduled.pop_front().expect("front checked");
1263 self.schedule_future_signal(
1264 signal,
1265 profile,
1266 representative_quote,
1267 &batch_quotes,
1268 &mut queued,
1269 &mut lifecycle,
1270 &mut future_executor,
1271 &mut portfolio,
1272 &pricer,
1273 &conversion_quotes,
1274 );
1275 processed_signals += 1;
1276 if should_report_progress(processed_signals, total_signals) {
1277 on_progress(ReplayProgress {
1278 processed_events,
1279 total_events,
1280 processed_signals,
1281 total_signals,
1282 });
1283 }
1284
1285 self.execute_queued_future(
1286 &batch_quotes,
1287 false,
1288 &mut queued,
1289 &mut lifecycle,
1290 &mut future_executor,
1291 &mut portfolio,
1292 &pricer,
1293 &conversion_quotes,
1294 );
1295 let increasing_symbols =
1296 queued_exposure_symbols(&queued, &batch_quotes, representative_quote.ts);
1297 invalid_quotes += self.settle_future_batch_symbols(
1298 &primary_quotes,
1299 &mut settled_quotes,
1300 Some(&increasing_symbols),
1301 &mut lifecycle,
1302 &mut future_executor,
1303 &mut portfolio,
1304 &pricer,
1305 &conversion_quotes,
1306 );
1307 self.execute_queued_future(
1308 &batch_quotes,
1309 true,
1310 &mut queued,
1311 &mut lifecycle,
1312 &mut future_executor,
1313 &mut portfolio,
1314 &pricer,
1315 &conversion_quotes,
1316 );
1317 }
1318
1319 invalid_quotes += self.settle_future_batch_symbols(
1321 &primary_quotes,
1322 &mut settled_quotes,
1323 None,
1324 &mut lifecycle,
1325 &mut future_executor,
1326 &mut portfolio,
1327 &pricer,
1328 &conversion_quotes,
1329 );
1330
1331 while scheduled
1332 .front()
1333 .is_some_and(|signal| signal.effective_ts == representative_quote.ts)
1334 {
1335 if is_cancelled() {
1336 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1337 }
1338 let signal = scheduled.pop_front().expect("front checked");
1339 self.schedule_future_signal(
1340 signal,
1341 profile,
1342 representative_quote,
1343 &batch_quotes,
1344 &mut queued,
1345 &mut lifecycle,
1346 &mut future_executor,
1347 &mut portfolio,
1348 &pricer,
1349 &conversion_quotes,
1350 );
1351 processed_signals += 1;
1352 if should_report_progress(processed_signals, total_signals) {
1353 on_progress(ReplayProgress {
1354 processed_events,
1355 total_events,
1356 processed_signals,
1357 total_signals,
1358 });
1359 }
1360 self.execute_queued_future(
1361 &batch_quotes,
1362 false,
1363 &mut queued,
1364 &mut lifecycle,
1365 &mut future_executor,
1366 &mut portfolio,
1367 &pricer,
1368 &conversion_quotes,
1369 );
1370 self.execute_queued_future(
1371 &batch_quotes,
1372 true,
1373 &mut queued,
1374 &mut lifecycle,
1375 &mut future_executor,
1376 &mut portfolio,
1377 &pricer,
1378 &conversion_quotes,
1379 );
1380 }
1381 }
1382
1383 for quote in &primary_quotes {
1384 observe_future_equity(
1385 &mut portfolio,
1386 &future_executor,
1387 quote.ts,
1388 &conversion_quotes,
1389 EquityObservationKind::PostOutput,
1390 &mut mtm_curve,
1391 &mut last_mtm_candidate,
1392 true,
1393 );
1394 processed_events += 1;
1395 if should_report_progress(processed_events, total_events) {
1396 on_progress(ReplayProgress {
1397 processed_events,
1398 total_events,
1399 processed_signals,
1400 total_signals,
1401 });
1402 }
1403 }
1404 if valuation_only {
1405 observe_future_equity(
1406 &mut portfolio,
1407 &future_executor,
1408 batch_ts,
1409 &conversion_quotes,
1410 EquityObservationKind::ConversionRevaluation,
1411 &mut mtm_curve,
1412 &mut last_mtm_candidate,
1413 false,
1414 );
1415 }
1416 conversion_quotes.retain_replay_causal_predecessors(
1417 batch_ts,
1418 scheduled
1419 .iter()
1420 .map(|signal| signal.effective_ts)
1421 .chain(primary_eod),
1422 );
1423
1424 if last_processed_primary_ts.is_some()
1425 && scheduled.is_empty()
1426 && queued.is_empty()
1427 && self.engine.open_positions().is_empty()
1428 && self.engine.pending_positions().is_empty()
1429 {
1430 effective_terminal_ts = last_processed_primary_ts;
1431 terminated_quiescently = true;
1432 break;
1433 }
1434 }
1435
1436 for action in queued {
1437 if is_cancelled() {
1438 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1439 }
1440 let mut disposition =
1441 ActionDisposition::rejected(action.action_id, "no_eligible_quote");
1442 disposition.action_kind = Some(action.action_kind);
1443 disposition.signal_ts = Some(action.signal_ts);
1444 disposition.effective_ts = Some(action.effective_ts);
1445 let _ = lifecycle.record(disposition);
1446 }
1447 for signal in scheduled {
1448 if is_cancelled() {
1449 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1450 }
1451 let mut disposition = ActionDisposition::rejected(
1452 format!("signal:{:08}", signal.sequence),
1453 "no_eligible_quote",
1454 );
1455 disposition.action_kind = Some(raw_signal_kind(&signal.signal).to_owned());
1456 disposition.signal_ts = Some(signal.signal_ts);
1457 disposition.effective_ts = Some(signal.effective_ts);
1458 let _ = lifecycle.record(disposition);
1459 processed_signals += 1;
1460 if should_report_progress(processed_signals, total_signals) {
1461 on_progress(ReplayProgress {
1462 processed_events,
1463 total_events,
1464 processed_signals,
1465 total_signals,
1466 });
1467 }
1468 }
1469
1470 if is_cancelled() {
1471 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1472 }
1473
1474 if self.config.close_on_finish {
1475 let execution_ts = effective_terminal_ts;
1476 let ids: Vec<String> = future_executor
1477 .open_snapshots()
1478 .into_iter()
1479 .map(|position| position.position_id)
1480 .collect();
1481 for (sequence, id) in ids.into_iter().enumerate() {
1482 if is_cancelled() {
1483 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1484 }
1485 let Some(symbol) = self
1486 .engine
1487 .get_position(&id)
1488 .map(|position| position.data.symbol.clone())
1489 else {
1490 continue;
1491 };
1492 let Some(quote) = portfolio.quote(&symbol).cloned() else {
1493 continue;
1494 };
1495 let action_id = format!("end_of_data:{sequence:08}");
1496 let transaction = (|| -> Result<_, FutureTransactionError> {
1497 let side = self
1498 .engine
1499 .get_position(&id)
1500 .ok_or_else(|| {
1501 FutureApplyError::Core(qs_core::CoreError::PositionNotFound(id.clone()))
1502 })?
1503 .data
1504 .side;
1505 let execution = pricer
1506 .market_exit(side, "e, self.pip_size(&symbol))
1507 .map_err(FutureApplyError::from)?;
1508 let engine_transaction =
1509 self.engine.begin_close_position_with_reason_future_at(
1510 &id,
1511 CloseReason::EndOfData,
1512 "e,
1513 execution,
1514 execution_ts.unwrap_or(quote.ts),
1515 )?;
1516 let affected =
1517 if FutureExecutor::requires_processing(engine_transaction.effects()) {
1518 match future_executor.process_future_effects_with_currency(
1519 engine_transaction.effects(),
1520 &self.engine,
1521 "e,
1522 Some(&action_id),
1523 None,
1524 execution_ts.unwrap_or(quote.ts),
1525 &mut portfolio,
1526 Some(&conversion_quotes),
1527 ) {
1528 Ok(affected) => affected,
1529 Err(error) => {
1530 engine_transaction.rollback(&mut self.engine);
1531 return Err(error.into());
1532 }
1533 }
1534 } else {
1535 Vec::new()
1536 };
1537 let _ = engine_transaction.commit();
1538 Ok(affected)
1539 })();
1540
1541 let mut disposition = match transaction {
1542 Ok(affected) => {
1543 let mut disposition = ActionDisposition::applied(action_id);
1544 disposition.position_ids = affected;
1545 disposition
1546 }
1547 Err(error) => ActionDisposition::failed(action_id, error.to_string()),
1548 };
1549 disposition.action_kind = Some("end_of_data".into());
1550 disposition.effective_ts = Some(execution_ts.unwrap_or(quote.ts));
1551 let _ = lifecycle.record(disposition);
1552 }
1553 }
1554
1555 if is_cancelled() {
1556 return Err(StreamingReplayError::Cancelled(ReplayCancelled));
1557 }
1558
1559 if let Some(ts) = effective_terminal_ts {
1560 let observation_kind = if terminated_quiescently {
1561 EquityObservationKind::QuiescentTermination
1562 } else {
1563 EquityObservationKind::EndOfData
1564 };
1565 observe_future_equity(
1566 &mut portfolio,
1567 &future_executor,
1568 ts,
1569 &conversion_quotes,
1570 observation_kind,
1571 &mut mtm_curve,
1572 &mut last_mtm_candidate,
1573 false,
1574 );
1575 future_executor.finalize_pending_orders_at_end(ts);
1576 }
1577
1578 let pending_orders = self
1579 .engine
1580 .pending_positions()
1581 .into_iter()
1582 .map(|position| {
1583 let metadata = future_executor.pending_metadata(&position.data.id);
1584 PendingOrderSnapshot {
1585 position_id: position.data.id.clone(),
1586 action_id: metadata.as_ref().map(|value| value.0.clone()),
1587 signal_ts: metadata.as_ref().map(|value| value.1),
1588 effective_ts: metadata.as_ref().map(|value| value.2),
1589 symbol: position.data.symbol.clone(),
1590 side: position.data.side,
1591 order_type: position.data.order_type,
1592 requested_price: position.data.pending_price,
1593 size: position.data.size,
1594 initial_stop: position.current_stoploss(),
1595 group: position.data.group.clone(),
1596 trade_id: position.data.trade_id.clone(),
1597 }
1598 })
1599 .collect();
1600 let mut tags = BTreeMap::new();
1601 tags.insert("invalid_quote_count".into(), invalid_quotes.to_string());
1602 tags.insert(
1603 "termination_reason".into(),
1604 if terminated_quiescently {
1605 "quiescent"
1606 } else {
1607 "end_of_data"
1608 }
1609 .into(),
1610 );
1611 insert_economic_support_metadata(&mut tags, &self.config);
1612 let (equity_curve, mtm_output_summary) = mtm_curve.into_parts();
1613 let artifacts = FutureBacktestArtifacts {
1614 format_version: FUTURE_ARTIFACT_FORMAT_VERSION,
1615 execution: ExecutionMetadata {
1616 execution_model,
1617 initial_balance: self.config.initial_balance,
1618 account_currency: future
1619 .currency_plan
1620 .as_ref()
1621 .map(|plan| plan.account_currency().to_owned()),
1622 currency_plan: future.currency_plan.clone(),
1623 contract_sizes: contract_sizes.into_iter().collect(),
1624 instrument_manifest: self.config.instrument_manifest.clone(),
1625 instrument_sizing: std::mem::take(&mut self.instrument_sizing),
1626 stale_quote_after_millis: future.stale_quote_after_ms,
1627 pnl_epsilon: future.pnl_epsilon,
1628 tags,
1629 ..ExecutionMetadata::default()
1630 },
1631 fills: future_executor.fills.clone(),
1632 close_events: future_executor.close_events.clone(),
1633 completed_positions: future_executor.completed_positions.clone(),
1634 open_positions: portfolio.latest_open_positions().to_vec(),
1635 pending_orders,
1636 pending_order_lifecycle: future_executor.pending_order_lifecycle,
1637 lifecycle,
1638 equity_curve,
1639 mtm_output_summary,
1640 max_drawdown: portfolio.max_drawdown(),
1641 max_drawdown_pct: portfolio.max_drawdown_pct(),
1642 };
1643 on_progress(ReplayProgress {
1644 processed_events,
1645 total_events: known_total_events.unwrap_or(processed_events),
1646 processed_signals,
1647 total_signals,
1648 });
1649 Ok(BacktestResult::from_future_artifacts_with_options(
1650 artifacts,
1651 self.evaluation_options,
1652 ))
1653 }
1654
1655 #[allow(clippy::too_many_arguments)]
1656 fn settle_future_batch_symbols(
1657 &mut self,
1658 primary_quotes: &[PriceQuote],
1659 settled_quotes: &mut [bool],
1660 symbols: Option<&BTreeSet<String>>,
1661 lifecycle: &mut LifecycleLedger,
1662 future_executor: &mut FutureExecutor,
1663 portfolio: &mut PortfolioRecorder,
1664 pricer: &ExecutionPricer,
1665 conversion_quotes: &ConversionQuoteBook,
1666 ) -> u64 {
1667 let mut failures = 0;
1668 for (index, quote) in primary_quotes.iter().enumerate() {
1669 if settled_quotes[index]
1670 || symbols.is_some_and(|symbols| !symbols.contains("e.symbol))
1671 {
1672 continue;
1673 }
1674 if self
1675 .settle_future_quote(
1676 quote,
1677 lifecycle,
1678 future_executor,
1679 portfolio,
1680 pricer,
1681 conversion_quotes,
1682 )
1683 .is_err()
1684 {
1685 failures += 1;
1686 }
1687 settled_quotes[index] = true;
1688 }
1689 failures
1690 }
1691
1692 #[allow(clippy::too_many_arguments)]
1693 fn settle_future_quote(
1694 &mut self,
1695 quote: &PriceQuote,
1696 lifecycle: &mut LifecycleLedger,
1697 future_executor: &mut FutureExecutor,
1698 portfolio: &mut PortfolioRecorder,
1699 pricer: &ExecutionPricer,
1700 conversion_quotes: &ConversionQuoteBook,
1701 ) -> Result<(), FutureTransactionError> {
1702 let (prepared, failures) = self.prepare_triggering_pending(quote, pricer);
1703 for (position_id, error) in failures {
1704 let action_id = format!("pending_execution:{position_id}");
1705 let mut disposition = ActionDisposition::rejected(action_id.clone(), error);
1706 disposition.action_kind = Some("pending_execution".into());
1707 disposition.effective_ts = Some(quote.ts);
1708 disposition.position_ids.push(position_id.clone());
1709 let _ = lifecycle.record(disposition);
1710 if let Ok(engine_transaction) = self.engine.begin_future_action(
1711 Action::CancelPending {
1712 position_id: position_id.clone(),
1713 },
1714 quote.ts,
1715 ) {
1716 if FutureExecutor::requires_processing(engine_transaction.effects())
1717 && let Err(error) = future_executor.process_future_effects_with_currency(
1718 engine_transaction.effects(),
1719 &self.engine,
1720 quote,
1721 Some(&action_id),
1722 None,
1723 quote.ts,
1724 portfolio,
1725 Some(conversion_quotes),
1726 )
1727 {
1728 engine_transaction.rollback(&mut self.engine);
1729 return Err(error.into());
1730 }
1731 let _ = engine_transaction.commit();
1732 }
1733 }
1734
1735 let pip_size = self.pip_size("e.symbol);
1736 let engine_transaction = self
1737 .engine
1738 .begin_on_price_future_effects_priced(quote, &prepared, pricer, pip_size)?;
1739 if FutureExecutor::requires_processing(engine_transaction.effects())
1740 && let Err(error) = future_executor.process_future_effects_with_currency(
1741 engine_transaction.effects(),
1742 &self.engine,
1743 quote,
1744 None,
1745 None,
1746 quote.ts,
1747 portfolio,
1748 Some(conversion_quotes),
1749 )
1750 {
1751 engine_transaction.rollback(&mut self.engine);
1752 return Err(error.into());
1753 }
1754 let _ = engine_transaction.commit();
1755 Ok(())
1756 }
1757
1758 #[allow(clippy::too_many_arguments)]
1759 fn schedule_future_signal(
1760 &mut self,
1761 scheduled: ScheduledSignal,
1762 profile: Option<&ManagementProfile>,
1763 quote: &PriceQuote,
1764 batch_quotes: &BTreeMap<String, PriceQuote>,
1765 queued: &mut VecDeque<QueuedAction>,
1766 lifecycle: &mut LifecycleLedger,
1767 future_executor: &mut FutureExecutor,
1768 portfolio: &mut PortfolioRecorder,
1769 pricer: &ExecutionPricer,
1770 conversion_quotes: &ConversionQuoteBook,
1771 ) {
1772 let base_id = format!("signal:{:08}", scheduled.sequence);
1773 if let RawSignal::Entry {
1774 symbol,
1775 side,
1776 order_type: OrderType::Market,
1777 ..
1778 } = &scheduled.signal
1779 {
1780 queued.push_back(QueuedAction {
1781 action_id: base_id,
1782 action_kind: "entry".into(),
1783 action: Action::Open {
1784 symbol: symbol.clone(),
1785 side: *side,
1786 order_type: OrderType::Market,
1787 price: None,
1788 size: 1.0,
1789 stoploss: None,
1790 targets: Vec::new(),
1791 rules: Vec::new(),
1792 group: None,
1793 trade_id: None,
1794 },
1795 execution: None,
1796 symbol: symbol.clone(),
1797 signal_ts: scheduled.signal_ts,
1798 effective_ts: scheduled.effective_ts,
1799 entry_signal: Some(scheduled.signal),
1800 entry_profile: profile.cloned(),
1801 });
1802 return;
1803 }
1804 if scheduled.signal.is_entry() {
1805 let resolved = match profile {
1806 Some(profile) => profile.apply_entry_signal(&scheduled.signal),
1807 None => resolve_unprofiled_entry(&scheduled.signal),
1808 };
1809 match resolved {
1810 Ok(Some(resolved)) => {
1811 let entry_quote = batch_quotes.get(&resolved.symbol).unwrap_or(quote);
1812 self.enqueue_resolved_entry(
1813 base_id,
1814 scheduled,
1815 resolved,
1816 entry_quote,
1817 lifecycle,
1818 future_executor,
1819 portfolio,
1820 pricer,
1821 conversion_quotes,
1822 )
1823 }
1824 Ok(None) => {
1825 let mut disposition = ActionDisposition::skipped(base_id, "not_an_entry");
1826 disposition.action_kind = Some("entry".into());
1827 disposition.signal_ts = Some(scheduled.signal_ts);
1828 disposition.effective_ts = Some(scheduled.effective_ts);
1829 let _ = lifecycle.record(disposition);
1830 }
1831 Err(error) => {
1832 let mut disposition = ActionDisposition::rejected(base_id, error.to_string());
1833 disposition.action_kind = Some("entry".into());
1834 disposition.signal_ts = Some(scheduled.signal_ts);
1835 disposition.effective_ts = Some(scheduled.effective_ts);
1836 let _ = lifecycle.record(disposition);
1837 }
1838 }
1839 return;
1840 }
1841
1842 let actions = self.resolve_future_actions(&scheduled.signal);
1843 if actions.is_empty() {
1844 let mut disposition = ActionDisposition::skipped(base_id, "position_not_found");
1845 disposition.action_kind = Some(raw_signal_kind(&scheduled.signal).to_owned());
1846 disposition.signal_ts = Some(scheduled.signal_ts);
1847 disposition.effective_ts = Some(scheduled.effective_ts);
1848 let _ = lifecycle.record(disposition);
1849 return;
1850 }
1851 for (index, action) in actions.into_iter().enumerate() {
1852 let action_id = format!("{base_id}:action:{index:03}");
1853 let Some(symbol) = self.action_symbol(&action) else {
1854 self.apply_future_action(
1855 action_id,
1856 raw_signal_kind(&scheduled.signal).to_owned(),
1857 action,
1858 None,
1859 scheduled.signal_ts,
1860 scheduled.effective_ts,
1861 quote,
1862 lifecycle,
1863 future_executor,
1864 portfolio,
1865 pricer,
1866 conversion_quotes,
1867 );
1868 continue;
1869 };
1870 if is_fill_bearing(&action) {
1871 queued.push_back(QueuedAction {
1872 action_id,
1873 action_kind: raw_signal_kind(&scheduled.signal).to_owned(),
1874 action,
1875 execution: None,
1876 symbol,
1877 signal_ts: scheduled.signal_ts,
1878 effective_ts: scheduled.effective_ts,
1879 entry_signal: None,
1880 entry_profile: None,
1881 });
1882 } else {
1883 let action_quote = batch_quotes.get(&symbol).unwrap_or(quote);
1884 self.apply_future_action(
1885 action_id,
1886 raw_signal_kind(&scheduled.signal).to_owned(),
1887 action,
1888 None,
1889 scheduled.signal_ts,
1890 scheduled.effective_ts,
1891 action_quote,
1892 lifecycle,
1893 future_executor,
1894 portfolio,
1895 pricer,
1896 conversion_quotes,
1897 );
1898 }
1899 }
1900 }
1901
1902 #[allow(clippy::too_many_arguments)]
1903 fn enqueue_resolved_entry(
1904 &mut self,
1905 action_id: String,
1906 scheduled: ScheduledSignal,
1907 resolved: ResolvedEntry,
1908 quote: &PriceQuote,
1909 lifecycle: &mut LifecycleLedger,
1910 future_executor: &mut FutureExecutor,
1911 portfolio: &mut PortfolioRecorder,
1912 pricer: &ExecutionPricer,
1913 conversion_quotes: &ConversionQuoteBook,
1914 ) {
1915 match self.finalize_resolved_entry(
1916 resolved,
1917 future_executor.balance(),
1918 scheduled.effective_ts,
1919 Some(conversion_quotes),
1920 ) {
1921 Ok(action) => self.apply_future_action(
1922 action_id,
1923 "entry".into(),
1924 action,
1925 None,
1926 scheduled.signal_ts,
1927 scheduled.effective_ts,
1928 quote,
1929 lifecycle,
1930 future_executor,
1931 portfolio,
1932 pricer,
1933 conversion_quotes,
1934 ),
1935 Err(error) => {
1936 let mut disposition = ActionDisposition::rejected(action_id, error);
1937 disposition.action_kind = Some("entry".into());
1938 disposition.signal_ts = Some(scheduled.signal_ts);
1939 disposition.effective_ts = Some(scheduled.effective_ts);
1940 let _ = lifecycle.record(disposition);
1941 }
1942 }
1943 }
1944
1945 #[allow(clippy::too_many_arguments)]
1946 fn execute_queued_future(
1947 &mut self,
1948 quotes: &BTreeMap<String, PriceQuote>,
1949 exposure_increasing: bool,
1950 queued: &mut VecDeque<QueuedAction>,
1951 lifecycle: &mut LifecycleLedger,
1952 future_executor: &mut FutureExecutor,
1953 portfolio: &mut PortfolioRecorder,
1954 pricer: &ExecutionPricer,
1955 conversion_quotes: &ConversionQuoteBook,
1956 ) {
1957 let mut remaining = VecDeque::new();
1958 while let Some(mut action) = queued.pop_front() {
1959 let increases = is_exposure_increasing(&action.action);
1960 let Some(quote) = quotes.get(&action.symbol) else {
1961 remaining.push_back(action);
1962 continue;
1963 };
1964 if action.effective_ts > quote.ts || increases != exposure_increasing {
1965 remaining.push_back(action);
1966 continue;
1967 }
1968
1969 if let Some(mut signal) = action.entry_signal.take() {
1970 let (side, symbol) = match &signal {
1971 RawSignal::Entry { side, symbol, .. } => (*side, symbol.clone()),
1972 _ => unreachable!("queued entry metadata must contain an entry signal"),
1973 };
1974 let execution = match pricer.market_entry(side, quote, self.pip_size(&symbol)) {
1975 Ok(fill) => fill,
1976 Err(error) => {
1977 let mut disposition =
1978 ActionDisposition::rejected(action.action_id, error.to_string());
1979 disposition.action_kind = Some(action.action_kind);
1980 disposition.signal_ts = Some(action.signal_ts);
1981 disposition.effective_ts = Some(action.effective_ts);
1982 let _ = lifecycle.record(disposition);
1983 continue;
1984 }
1985 };
1986 if let RawSignal::Entry { price, .. } = &mut signal {
1987 *price = Some(execution.price);
1988 }
1989 action.execution = Some(execution);
1990 let resolved = match action.entry_profile.as_ref() {
1991 Some(profile) => profile.apply_entry_signal(&signal),
1992 None => resolve_unprofiled_entry(&signal),
1993 };
1994 match resolved {
1995 Ok(Some(resolved)) => match self.finalize_resolved_entry(
1996 resolved,
1997 future_executor.balance(),
1998 quote.ts,
1999 Some(conversion_quotes),
2000 ) {
2001 Ok(finalized) => action.action = finalized,
2002 Err(error) => {
2003 let mut disposition =
2004 ActionDisposition::rejected(action.action_id, error);
2005 disposition.action_kind = Some(action.action_kind);
2006 disposition.signal_ts = Some(action.signal_ts);
2007 disposition.effective_ts = Some(action.effective_ts);
2008 let _ = lifecycle.record(disposition);
2009 continue;
2010 }
2011 },
2012 Ok(None) => {
2013 let mut disposition =
2014 ActionDisposition::skipped(action.action_id, "not_an_entry");
2015 disposition.action_kind = Some(action.action_kind);
2016 disposition.signal_ts = Some(action.signal_ts);
2017 disposition.effective_ts = Some(action.effective_ts);
2018 let _ = lifecycle.record(disposition);
2019 continue;
2020 }
2021 Err(error) => {
2022 let mut disposition =
2023 ActionDisposition::rejected(action.action_id, error.to_string());
2024 disposition.action_kind = Some(action.action_kind);
2025 disposition.signal_ts = Some(action.signal_ts);
2026 disposition.effective_ts = Some(action.effective_ts);
2027 let _ = lifecycle.record(disposition);
2028 continue;
2029 }
2030 }
2031 }
2032 self.apply_future_action(
2033 action.action_id,
2034 action.action_kind,
2035 action.action,
2036 action.execution,
2037 action.signal_ts,
2038 action.effective_ts,
2039 quote,
2040 lifecycle,
2041 future_executor,
2042 portfolio,
2043 pricer,
2044 conversion_quotes,
2045 );
2046 }
2047 *queued = remaining;
2048 }
2049
2050 #[allow(clippy::too_many_arguments)]
2051 fn apply_future_action(
2052 &mut self,
2053 action_id: String,
2054 action_kind: String,
2055 mut action: Action,
2056 execution: Option<ExecutionFill>,
2057 signal_ts: NaiveDateTime,
2058 effective_ts: NaiveDateTime,
2059 quote: &PriceQuote,
2060 lifecycle: &mut LifecycleLedger,
2061 future_executor: &mut FutureExecutor,
2062 portfolio: &mut PortfolioRecorder,
2063 pricer: &ExecutionPricer,
2064 conversion_quotes: &ConversionQuoteBook,
2065 ) {
2066 if let Action::Open {
2067 trade_id: Some(trade_id),
2068 ..
2069 } = &action
2070 && self.engine.manager.id_by_trade_id(trade_id).is_some()
2071 {
2072 let mut disposition = ActionDisposition::rejected(action_id, "duplicate_trade_id");
2073 disposition.action_kind = Some(action_kind);
2074 disposition.signal_ts = Some(signal_ts);
2075 disposition.effective_ts = Some(effective_ts);
2076 let _ = lifecycle.record(disposition);
2077 return;
2078 }
2079 if let Action::ScaleIn { position_id, .. } = &action
2080 && future_executor.has_close(position_id)
2081 {
2082 let mut disposition =
2083 ActionDisposition::rejected(action_id, "scale_in_after_close_not_supported");
2084 disposition.action_kind = Some(action_kind);
2085 disposition.signal_ts = Some(signal_ts);
2086 disposition.effective_ts = Some(effective_ts);
2087 disposition.position_ids.push(position_id.clone());
2088 let _ = lifecycle.record(disposition);
2089 return;
2090 }
2091
2092 let execution = match self.prepare_future_action(&mut action, execution, quote, pricer) {
2093 Ok(execution) => execution,
2094 Err(reason) => {
2095 let mut disposition = ActionDisposition::rejected(action_id, reason);
2096 disposition.action_kind = Some(action_kind);
2097 disposition.signal_ts = Some(signal_ts);
2098 disposition.effective_ts = Some(effective_ts);
2099 let _ = lifecycle.record(disposition);
2100 return;
2101 }
2102 };
2103
2104 let engine_transaction = match execution {
2105 Some(execution) => self
2106 .engine
2107 .begin_priced_future_action(action, quote, execution),
2108 None => self.engine.begin_future_action(action, effective_ts),
2109 };
2110 let engine_transaction = match engine_transaction {
2111 Ok(transaction) => transaction,
2112 Err(error) => {
2113 let mut disposition = ActionDisposition::rejected(action_id, error.to_string());
2114 disposition.action_kind = Some(action_kind);
2115 disposition.signal_ts = Some(signal_ts);
2116 disposition.effective_ts = Some(effective_ts);
2117 let _ = lifecycle.record(disposition);
2118 return;
2119 }
2120 };
2121
2122 let mut affected = if FutureExecutor::requires_processing(engine_transaction.effects()) {
2123 match future_executor.process_future_effects_with_currency(
2124 engine_transaction.effects(),
2125 &self.engine,
2126 quote,
2127 Some(&action_id),
2128 Some(signal_ts),
2129 effective_ts,
2130 portfolio,
2131 Some(conversion_quotes),
2132 ) {
2133 Ok(affected) => affected,
2134 Err(error) => {
2135 engine_transaction.rollback(&mut self.engine);
2136 let mut disposition = ActionDisposition::failed(action_id, error.to_string());
2137 disposition.action_kind = Some(action_kind);
2138 disposition.signal_ts = Some(signal_ts);
2139 disposition.effective_ts = Some(effective_ts);
2140 let _ = lifecycle.record(disposition);
2141 return;
2142 }
2143 }
2144 } else {
2145 Vec::new()
2146 };
2147 for future_effect in engine_transaction.effects() {
2148 match future_effect.effect() {
2149 Effect::OrderPlaced { id } | Effect::OrderCancelled { id } => {
2150 affected.push(id.clone());
2151 }
2152 _ => {}
2153 }
2154 }
2155 affected.sort();
2156 affected.dedup();
2157 let _ = engine_transaction.commit();
2158
2159 let mut disposition = ActionDisposition::applied(action_id);
2160 disposition.action_kind = Some(action_kind);
2161 disposition.signal_ts = Some(signal_ts);
2162 disposition.effective_ts = Some(effective_ts);
2163 disposition.position_ids = affected;
2164 let _ = lifecycle.record(disposition);
2165 }
2166
2167 fn prepare_future_action(
2168 &self,
2169 action: &mut Action,
2170 prepriced: Option<ExecutionFill>,
2171 quote: &PriceQuote,
2172 pricer: &ExecutionPricer,
2173 ) -> Result<Option<ExecutionFill>, String> {
2174 let mut execution = prepriced;
2175 match action {
2176 Action::Open {
2177 symbol,
2178 side,
2179 order_type,
2180 price,
2181 size,
2182 ..
2183 } => {
2184 if !valid_accounting_size(*size) {
2185 return Err(format!(
2186 "position size must be finite and greater than the accounting tolerance, got {size}"
2187 ));
2188 }
2189 if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
2190 return Err(format!(
2191 "supplied entry price must be finite and positive, got {price:?}"
2192 ));
2193 }
2194 if *order_type == OrderType::Market {
2195 let priced = match execution {
2196 Some(priced) => priced,
2197 None => pricer
2198 .market_entry(*side, quote, self.pip_size(symbol))
2199 .map_err(|error| error.to_string())?,
2200 };
2201 *price = Some(priced.price);
2202 execution = Some(priced);
2203 } else {
2204 if price.is_none() {
2205 return Err("pending entry requires a requested price".to_owned());
2206 }
2207 execution = None;
2208 }
2209 }
2210 Action::ScaleIn {
2211 position_id,
2212 price,
2213 size,
2214 ..
2215 } => {
2216 if !valid_accounting_size(*size) {
2217 return Err(format!(
2218 "scale-in size must be finite and greater than the accounting tolerance, got {size}"
2219 ));
2220 }
2221 if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
2222 return Err(format!(
2223 "supplied scale-in price must be finite and positive, got {price:?}"
2224 ));
2225 }
2226 let side = self
2227 .engine
2228 .get_position(position_id)
2229 .map(|position| position.data.side)
2230 .ok_or_else(|| format!("position not found: {position_id}"))?;
2231 let priced = match execution {
2232 Some(priced) => priced,
2233 None => pricer
2234 .market_entry(side, quote, self.pip_size("e.symbol))
2235 .map_err(|error| error.to_string())?,
2236 };
2237 *price = Some(priced.price);
2238 execution = Some(priced);
2239 }
2240 Action::ClosePosition { position_id } | Action::ClosePartial { position_id, .. } => {
2241 let position = self
2242 .engine
2243 .get_position(position_id)
2244 .ok_or_else(|| format!("position not found: {position_id}"))?;
2245 if position.data.symbol != quote.symbol {
2246 return Err(format!(
2247 "position symbol {} does not match quote symbol {}",
2248 position.data.symbol, quote.symbol
2249 ));
2250 }
2251 execution = Some(
2252 pricer
2253 .market_exit(
2254 position.data.side,
2255 quote,
2256 self.pip_size(&position.data.symbol),
2257 )
2258 .map_err(|error| error.to_string())?,
2259 );
2260 }
2261 _ => execution = None,
2262 }
2263 Ok(execution)
2264 }
2265
2266 fn prepare_triggering_pending(
2267 &self,
2268 quote: &PriceQuote,
2269 pricer: &ExecutionPricer,
2270 ) -> (Vec<PreparedPendingFill>, Vec<(String, String)>) {
2271 let ids = self
2272 .engine
2273 .manager
2274 .pending_ids_by_symbol_sorted("e.symbol);
2275 let mut prepared = Vec::new();
2276 let mut failures = Vec::new();
2277 for id in ids {
2278 let Some(position) = self.engine.get_position(&id) else {
2279 continue;
2280 };
2281 let Some(purpose) = position.pending_fill_purpose(quote, self.config.fill_model) else {
2282 continue;
2283 };
2284 let execution = match pricer.price(
2285 purpose,
2286 position.data.side,
2287 quote,
2288 position.data.pending_price,
2289 self.pip_size("e.symbol),
2290 ) {
2291 Ok(fill) => fill,
2292 Err(error) => {
2293 failures.push((id, error.to_string()));
2294 continue;
2295 }
2296 };
2297
2298 let size = position.data.size;
2299 if !valid_accounting_size(size) {
2300 failures.push((
2301 id,
2302 format!(
2303 "pending size must be finite and greater than the accounting tolerance, got {size}"
2304 ),
2305 ));
2306 continue;
2307 }
2308
2309 prepared.push(PreparedPendingFill {
2310 position_id: id,
2311 execution,
2312 size,
2313 });
2314 }
2315 (prepared, failures)
2316 }
2317
2318 fn pip_size(&self, symbol: &str) -> f64 {
2319 self.config
2320 .symbol_specs
2321 .get(symbol)
2322 .map(|spec| 10_f64.powi(-(spec.pip_position as i32)))
2323 .unwrap_or(0.0001)
2324 }
2325
2326 fn action_symbol(&self, action: &Action) -> Option<String> {
2327 match action {
2328 Action::Open { symbol, .. } => Some(symbol.clone()),
2329 Action::ClosePosition { position_id }
2330 | Action::ClosePartial { position_id, .. }
2331 | Action::ModifyStoploss { position_id, .. }
2332 | Action::MoveStoplossToEntry { position_id }
2333 | Action::AddTarget { position_id, .. }
2334 | Action::RemoveTarget { position_id, .. }
2335 | Action::ModifyTarget { position_id, .. }
2336 | Action::AddRule { position_id, .. }
2337 | Action::RemoveRule { position_id, .. }
2338 | Action::ScaleIn { position_id, .. }
2339 | Action::CancelPending { position_id } => self
2340 .engine
2341 .get_position(position_id)
2342 .map(|position| position.data.symbol.clone()),
2343 Action::CloseAllOf { symbol } | Action::ModifyAllStoploss { symbol, .. } => {
2344 Some(symbol.clone())
2345 }
2346 _ => None,
2347 }
2348 }
2349
2350 fn resolve_future_actions(&self, signal: &RawSignal) -> Vec<Action> {
2351 match signal {
2352 RawSignal::CloseAllOf { symbol, .. } => self
2353 .engine
2354 .manager
2355 .open_ids_by_symbol_sorted(symbol)
2356 .into_iter()
2357 .map(|position_id| Action::ClosePosition { position_id })
2358 .collect(),
2359 RawSignal::CloseAll { .. } => self
2360 .engine
2361 .manager
2362 .ids_by_status_sorted(PositionStatus::Open)
2363 .into_iter()
2364 .map(|position_id| Action::ClosePosition { position_id })
2365 .collect(),
2366 RawSignal::CloseAllInGroup { group_id, .. } => {
2367 let mut ids = self.engine.manager.open_ids_by_group(group_id);
2368 ids.sort();
2369 ids.into_iter()
2370 .map(|position_id| Action::ClosePosition { position_id })
2371 .collect()
2372 }
2373 RawSignal::CancelAllPending { .. } => self
2374 .engine
2375 .manager
2376 .ids_by_status_sorted(PositionStatus::Pending)
2377 .into_iter()
2378 .map(|position_id| Action::CancelPending { position_id })
2379 .collect(),
2380 _ => resolve_signal(signal, &self.engine),
2381 }
2382 }
2383}
2384
2385#[allow(clippy::too_many_arguments)]
2386fn observe_future_equity(
2387 portfolio: &mut PortfolioRecorder,
2388 future_executor: &FutureExecutor,
2389 ts: NaiveDateTime,
2390 conversion_quotes: &ConversionQuoteBook,
2391 kind: EquityObservationKind,
2392 collector: &mut MtmCurveCollector,
2393 last_candidate: &mut Option<EquityPoint>,
2394 suppress_unchanged_post_output: bool,
2395) {
2396 portfolio.set_realized_pnl(future_executor.realized_pnl());
2397 let mut point = portfolio.observe_with_currency(
2398 ts,
2399 future_executor.open_snapshots(),
2400 Some(conversion_quotes),
2401 );
2402 point.observation_kind = Some(kind.as_str().to_owned());
2403 if suppress_unchanged_post_output
2404 && last_candidate
2405 .as_ref()
2406 .is_some_and(|previous| same_equity_values(previous, &point))
2407 {
2408 return;
2409 }
2410 collector.observe(point.clone());
2411 *last_candidate = Some(point);
2412}
2413
2414fn same_equity_values(left: &EquityPoint, right: &EquityPoint) -> bool {
2415 let mut left = left.clone();
2416 let mut right = right.clone();
2417 left.observation_kind = None;
2418 left.observation_sequence = None;
2419 right.observation_kind = None;
2420 right.observation_sequence = None;
2421 left == right
2422}
2423
2424fn valid_accounting_size(size: f64) -> bool {
2425 size.is_finite() && size > position_size_tolerance(size)
2426}
2427
2428fn explicit_instrument_spec<'a>(
2429 config: &'a BacktestConfig,
2430 symbol: &str,
2431) -> Option<&'a InstrumentSpec> {
2432 config
2433 .instrument_manifest
2434 .as_ref()?
2435 .instruments
2436 .get(symbol)
2437 .map(|artifact| &artifact.spec)
2438}
2439
2440fn decimal_to_f64(value: Decimal, field: &str) -> Result<f64, String> {
2441 let value = value
2442 .to_string()
2443 .parse::<f64>()
2444 .map_err(|error| format!("invalid {field}: {error}"))?;
2445 if value.is_finite() {
2446 Ok(value)
2447 } else {
2448 Err(format!("{field} must be finite"))
2449 }
2450}
2451
2452fn instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
2453 decimal_to_f64(
2454 spec.economics.contract_multiplier.get(),
2455 "instrument contract multiplier",
2456 )
2457 .and_then(|value| {
2458 if value > 0.0 {
2459 Ok(value)
2460 } else {
2461 Err("instrument contract multiplier must be positive".into())
2462 }
2463 })
2464}
2465
2466fn supported_instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
2467 if spec.status != ListingStatus::Trading {
2468 return Err(format!(
2469 "instrument {} is not in trading status",
2470 spec.instrument
2471 ));
2472 }
2473 if spec.economics.quantity_unit != QuantityUnit::StandardLot {
2474 return Err(format!(
2475 "unsupported quantity unit for instrument {}: {:?}",
2476 spec.instrument, spec.economics.quantity_unit
2477 ));
2478 }
2479 let model = spec.economics.pnl_model.as_str();
2480 if model != EconomicsModelId::FX_QUOTE_LINEAR_V1
2481 && model != EconomicsModelId::CFD_QUOTE_LINEAR_V1
2482 {
2483 return Err(format!(
2484 "unsupported P&L model for instrument {}: {model}",
2485 spec.instrument
2486 ));
2487 }
2488 instrument_multiplier(spec)
2489}
2490
2491fn validate_instrument_manifest(config: &BacktestConfig) -> Result<(), String> {
2492 let Some(manifest) = &config.instrument_manifest else {
2493 return Ok(());
2494 };
2495 for (symbol, artifact) in &manifest.instruments {
2496 if symbol.is_empty() {
2497 return Err("instrument manifest symbol must not be empty".into());
2498 }
2499 artifact
2500 .spec
2501 .validate()
2502 .map_err(|error| format!("invalid instrument spec for {symbol}: {error}"))?;
2503 if artifact.resolved.instrument != artifact.spec.instrument {
2504 return Err(format!(
2505 "resolved instrument and spec identity differ for {symbol}"
2506 ));
2507 }
2508 if artifact.resolved.spec_revision != artifact.spec.revision {
2509 return Err(format!(
2510 "resolved specification revision does not match the embedded spec for {symbol}"
2511 ));
2512 }
2513 supported_instrument_multiplier(&artifact.spec)?;
2514 }
2515 for binding in &manifest.stored_series {
2516 let known = manifest
2517 .instruments
2518 .values()
2519 .any(|artifact| artifact.resolved == binding.instrument);
2520 if !known {
2521 return Err(format!(
2522 "stored series {}:{} references an instrument outside the manifest",
2523 binding.source_partition, binding.source_symbol
2524 ));
2525 }
2526 let artifact = manifest
2527 .instruments
2528 .values()
2529 .find(|artifact| artifact.resolved == binding.instrument)
2530 .expect("known binding reference has an instrument artifact");
2531 if binding.effective != artifact.spec.effective {
2532 return Err(format!(
2533 "stored series {}:{} effective interval differs from its instrument spec",
2534 binding.source_partition, binding.source_symbol
2535 ));
2536 }
2537 }
2538 Ok(())
2539}
2540
2541fn effective_contract_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
2542 let mut contract_sizes = config.contract_sizes.clone();
2543 if let Some(manifest) = &config.instrument_manifest {
2544 for (symbol, artifact) in &manifest.instruments {
2545 if let Ok(multiplier) = instrument_multiplier(&artifact.spec) {
2546 contract_sizes.insert(symbol.clone(), multiplier);
2547 }
2548 }
2549 }
2550 contract_sizes
2551}
2552
2553fn is_monetary_sizing(policy: &SizingPolicy) -> bool {
2554 matches!(
2555 policy,
2556 SizingPolicy::FixedRiskAmount { .. } | SizingPolicy::BalanceRiskPercent { .. }
2557 )
2558}
2559
2560fn accept_legacy_quote(
2561 quote: &PriceQuote,
2562 last_quote_ts: &mut BTreeMap<String, NaiveDateTime>,
2563) -> bool {
2564 if ExecutionPricer::validate_quote(quote).is_err()
2565 || last_quote_ts
2566 .get("e.symbol)
2567 .is_some_and(|last| *last > quote.ts)
2568 {
2569 return false;
2570 }
2571 last_quote_ts.insert(quote.symbol.clone(), quote.ts);
2572 true
2573}
2574
2575fn validate_replay_config(
2576 config: &BacktestConfig,
2577 future: Option<&FutureQuoteConfig>,
2578 raw_signals: &[RawSignal],
2579) -> Result<(), String> {
2580 if !config.initial_balance.is_finite() || config.initial_balance <= 0.0 {
2581 return Err(format!(
2582 "initial balance must be finite and positive, got {}",
2583 config.initial_balance
2584 ));
2585 }
2586 for (symbol, contract_size) in &config.contract_sizes {
2587 if symbol.is_empty() {
2588 return Err("contract-size symbol must not be empty".into());
2589 }
2590 if !contract_size.is_finite() || *contract_size <= 0.0 {
2591 return Err(format!(
2592 "contract size for {symbol} must be finite and positive, got {contract_size}"
2593 ));
2594 }
2595 }
2596
2597 validate_instrument_manifest(config)?;
2598 for (symbol, spec) in &config.symbol_specs {
2599 validate_symbol_spec(symbol, spec)?;
2600 if explicit_instrument_spec(config, symbol).is_none() {
2601 resolve_legacy_economics(spec).map_err(|error| error.to_string())?;
2602 }
2603 }
2604
2605 let entry_symbols: Vec<&str> = raw_signals
2606 .iter()
2607 .filter_map(|signal| match signal {
2608 RawSignal::Entry { symbol, .. } => Some(symbol.as_str()),
2609 _ => None,
2610 })
2611 .collect();
2612 if !entry_symbols.is_empty() && config.sizing.is_none() {
2613 return Err("raw entry requires BacktestConfig.sizing".to_owned());
2614 }
2615 if let Some(policy) = &config.sizing {
2616 validate_sizing_policy(policy)?;
2617 for symbol in &entry_symbols {
2618 if !config.symbol_specs.contains_key(*symbol)
2619 && explicit_instrument_spec(config, symbol).is_none()
2620 {
2621 return Err(format!("missing instrument or symbol spec for {symbol}"));
2622 }
2623 }
2624 if !entry_symbols.is_empty() && is_monetary_sizing(policy) {
2625 let future = future.ok_or_else(|| {
2626 "monetary sizing requires FutureQuote execution and a currency plan".to_owned()
2627 })?;
2628 let plan = future
2629 .currency_plan
2630 .as_ref()
2631 .ok_or_else(|| "monetary sizing requires a FutureQuote currency plan".to_owned())?;
2632 for symbol in &entry_symbols {
2633 if plan.route_for_primary_symbol(symbol).is_none() {
2634 return Err(format!(
2635 "currency plan has no frozen route for primary symbol {symbol}"
2636 ));
2637 }
2638 }
2639 }
2640 }
2641
2642 if let Some(future) = future {
2643 if future.signal_latency_ms < 0 {
2644 return Err(format!(
2645 "signal latency must be non-negative, got {}",
2646 future.signal_latency_ms
2647 ));
2648 }
2649 let latency = Duration::milliseconds(future.signal_latency_ms);
2650 for signal in raw_signals {
2651 if signal.ts().checked_add_signed(latency).is_none() {
2652 return Err(format!(
2653 "signal latency overflows datetime for signal at {}",
2654 signal.ts()
2655 ));
2656 }
2657 }
2658 if !future.slippage_pips.is_finite() {
2659 return Err(format!(
2660 "slippage pips must be finite, got {}",
2661 future.slippage_pips
2662 ));
2663 }
2664 if future.stale_quote_after_ms.is_some_and(|value| value < 0) {
2665 return Err("stale quote threshold must be non-negative".into());
2666 }
2667 if !future.pnl_epsilon.is_finite() || future.pnl_epsilon < 0.0 {
2668 return Err(format!(
2669 "P&L epsilon must be finite and non-negative, got {}",
2670 future.pnl_epsilon
2671 ));
2672 }
2673 if future.conversion_stale_after_ms < 0 {
2674 return Err("conversion quote threshold must be non-negative".to_owned());
2675 }
2676 future
2677 .mtm_output
2678 .validate()
2679 .map_err(|error| error.to_string())?;
2680 }
2681 Ok(())
2682}
2683
2684fn validate_sizing_policy(policy: &SizingPolicy) -> Result<(), String> {
2685 let (name, value) = match policy {
2686 SizingPolicy::FixedLot { lots } => ("fixed lots", *lots),
2687 SizingPolicy::FixedRiskAmount { amount } => ("fixed risk amount", *amount),
2688 SizingPolicy::BalanceRiskPercent { percent } => ("balance risk percent", *percent),
2689 };
2690 if value.is_finite() && value > 0.0 {
2691 Ok(())
2692 } else {
2693 Err(format!("{name} must be finite and positive, got {value}"))
2694 }
2695}
2696
2697fn validate_symbol_spec(symbol: &str, spec: &qs_symbols::SymbolSpec) -> Result<(), String> {
2698 if symbol.is_empty() || spec.canonical.is_empty() {
2699 return Err("symbol spec names must not be empty".into());
2700 }
2701 if spec.digits > 18 || spec.pip_position > spec.digits {
2702 return Err(format!(
2703 "invalid price precision for {symbol}: digits={}, pip_position={}",
2704 spec.digits, spec.pip_position
2705 ));
2706 }
2707 if spec.lot_base_units <= 0
2708 || spec.lot_step_units <= 0
2709 || spec.lot_min_steps <= 0
2710 || spec.lot_max_steps < 0
2711 || (spec.lot_max_steps > 0 && spec.lot_max_steps < spec.lot_min_steps)
2712 {
2713 return Err(format!("invalid lot metadata for {symbol}"));
2714 }
2715 let lot_step = spec.lot_step();
2716 let min_lot = spec.lot_min();
2717 let max_lot = spec.lot_max();
2718 if !lot_step.is_finite()
2719 || lot_step <= 0.0
2720 || !min_lot.is_finite()
2721 || min_lot <= 0.0
2722 || !max_lot.is_finite()
2723 {
2724 return Err(format!("invalid derived lot metadata for {symbol}"));
2725 }
2726 Ok(())
2727}
2728
2729fn rejected_legacy_result(config: &BacktestConfig) -> BacktestResult {
2730 BacktestResult::from_trade_log(
2731 if config.initial_balance.is_finite() {
2732 config.initial_balance
2733 } else {
2734 0.0
2735 },
2736 Vec::new(),
2737 )
2738}
2739
2740fn rejected_future_result(
2741 config: &BacktestConfig,
2742 future: &FutureQuoteConfig,
2743 evaluation_options: EvaluationOptions,
2744 error: String,
2745) -> BacktestResult {
2746 let execution_model = ExecutionModel::new(
2747 qs_core::types::ExecutionConvention::FutureQuoteV1,
2748 config.fill_model,
2749 if future.slippage_pips == 0.0 {
2750 SlippageModel::None
2751 } else {
2752 SlippageModel::FixedPips {
2753 pips: future.slippage_pips,
2754 }
2755 },
2756 );
2757 let mut lifecycle = LifecycleLedger::new();
2758 let _ = lifecycle.record(ActionDisposition::rejected(
2759 "configuration",
2760 format!("invalid_configuration: {error}"),
2761 ));
2762 let mut tags = BTreeMap::new();
2763 tags.insert("configuration_error".into(), error);
2764 insert_economic_support_metadata(&mut tags, config);
2765 let artifacts = FutureBacktestArtifacts {
2766 execution: ExecutionMetadata {
2767 execution_model,
2768 initial_balance: if config.initial_balance.is_finite() {
2769 config.initial_balance
2770 } else {
2771 0.0
2772 },
2773 account_currency: future
2774 .currency_plan
2775 .as_ref()
2776 .map(|plan| plan.account_currency().to_owned()),
2777 currency_plan: future.currency_plan.clone(),
2778 contract_sizes: effective_contract_sizes(config)
2779 .into_iter()
2780 .filter(|(symbol, size)| !symbol.is_empty() && size.is_finite() && *size > 0.0)
2781 .collect(),
2782 instrument_manifest: config.instrument_manifest.clone(),
2783 instrument_sizing: Vec::new(),
2784 stale_quote_after_millis: future.stale_quote_after_ms,
2785 pnl_epsilon: if future.pnl_epsilon.is_finite() && future.pnl_epsilon >= 0.0 {
2786 future.pnl_epsilon
2787 } else {
2788 crate::artifacts::DEFAULT_PNL_EPSILON
2789 },
2790 tags,
2791 ..ExecutionMetadata::default()
2792 },
2793 lifecycle,
2794 mtm_output_summary: MtmOutputSummary {
2795 policy: future.mtm_output,
2796 ..MtmOutputSummary::default()
2797 },
2798 ..FutureBacktestArtifacts::default()
2799 };
2800 BacktestResult::from_future_artifacts_with_options(artifacts, evaluation_options)
2801}
2802
2803fn insert_economic_support_metadata(tags: &mut BTreeMap<String, String>, config: &BacktestConfig) {
2804 let mut compatibility_specs = config.symbol_specs.iter().peekable();
2805 if compatibility_specs.peek().is_none() {
2806 return;
2807 }
2808 tags.insert(
2809 "economics.guard".into(),
2810 LEGACY_ECONOMIC_GUARD_ID.to_owned(),
2811 );
2812 for (symbol, spec) in compatibility_specs {
2813 let prefix = format!("economics.symbol.{symbol}");
2814 tags.insert(format!("{prefix}.category"), spec.category.clone());
2815 match resolve_legacy_economics(spec) {
2816 Ok(economics) => {
2817 tags.insert(format!("{prefix}.status"), "supported".into());
2818 tags.insert(format!("{prefix}.model"), economics.model.as_str().into());
2819 tags.insert(
2820 format!("{prefix}.contract_multiplier"),
2821 economics.contract_multiplier.to_string(),
2822 );
2823 }
2824 Err(error) => {
2825 tags.insert(format!("{prefix}.status"), "unsupported".into());
2826 tags.insert(format!("{prefix}.reason"), error.to_string());
2827 }
2828 }
2829 }
2830}
2831
2832fn queued_exposure_symbols(
2833 queued: &VecDeque<QueuedAction>,
2834 quotes: &BTreeMap<String, PriceQuote>,
2835 batch_ts: NaiveDateTime,
2836) -> BTreeSet<String> {
2837 queued
2838 .iter()
2839 .filter(|action| {
2840 action.effective_ts <= batch_ts
2841 && quotes.contains_key(&action.symbol)
2842 && is_exposure_increasing(&action.action)
2843 })
2844 .map(|action| action.symbol.clone())
2845 .collect()
2846}
2847
2848fn is_exposure_increasing(action: &Action) -> bool {
2849 matches!(
2850 action,
2851 Action::Open {
2852 order_type: OrderType::Market,
2853 ..
2854 } | Action::ScaleIn { .. }
2855 )
2856}
2857
2858fn is_fill_bearing(action: &Action) -> bool {
2859 matches!(
2860 action,
2861 Action::Open {
2862 order_type: OrderType::Market,
2863 ..
2864 } | Action::ClosePosition { .. }
2865 | Action::ClosePartial { .. }
2866 | Action::ScaleIn { .. }
2867 )
2868}
2869
2870fn raw_signal_kind(signal: &RawSignal) -> &'static str {
2871 match signal {
2872 RawSignal::Entry { .. } => "entry",
2873 RawSignal::Close { .. } => "close",
2874 RawSignal::ClosePartial { .. } => "close_partial",
2875 RawSignal::ModifyStoploss { .. } => "modify_stoploss",
2876 RawSignal::MoveStoplossToEntry { .. } => "move_stoploss_to_entry",
2877 RawSignal::AddTarget { .. } => "add_target",
2878 RawSignal::RemoveTarget { .. } => "remove_target",
2879 RawSignal::ModifyTarget { .. } => "modify_target",
2880 RawSignal::AddRule { .. } => "add_rule",
2881 RawSignal::RemoveRule { .. } => "remove_rule",
2882 RawSignal::ScaleIn { .. } => "scale_in",
2883 RawSignal::CancelPending { .. } => "cancel_pending",
2884 RawSignal::CloseAllOf { .. } => "close_all_of",
2885 RawSignal::CloseAll { .. } => "close_all",
2886 RawSignal::CancelAllPending { .. } => "cancel_all_pending",
2887 RawSignal::ModifyAllStoploss { .. } => "modify_all_stoploss",
2888 RawSignal::CloseAllInGroup { .. } => "close_all_in_group",
2889 RawSignal::ModifyAllStoplossInGroup { .. } => "modify_all_stoploss_in_group",
2890 }
2891}
2892
2893#[cfg(test)]
2896mod tests {
2897 use super::*;
2898 use crate::currency::{ConversionRoute, FxPair};
2899 use crate::data_feed::{EventMetadata, FeedEvent, MarketEvent, SeriesRoles, VecFeed};
2900 use crate::profile::{ManagementProfile, PositionRef, RawSignal, StoplossMode};
2901 use chrono::NaiveDate;
2902 use qs_core::types::{CloseReason, FillPurpose, OrderType, Side, TargetSpec};
2903
2904 fn ts(h: u32, m: u32, s: u32) -> chrono::NaiveDateTime {
2905 NaiveDate::from_ymd_opt(2026, 1, 1)
2906 .unwrap()
2907 .and_hms_opt(h, m, s)
2908 .unwrap()
2909 }
2910
2911 fn tick(symbol: &str, bid: f64, ask: f64, time: chrono::NaiveDateTime) -> MarketEvent {
2912 MarketEvent::Tick {
2913 symbol: symbol.into(),
2914 ts: time,
2915 bid,
2916 ask,
2917 }
2918 }
2919
2920 fn test_symbol_spec(symbol: &str) -> qs_symbols::SymbolSpec {
2921 qs_symbols::SymbolSpec {
2922 canonical: symbol.to_ascii_lowercase(),
2923 pip_position: 4,
2924 digits: 5,
2925 category: "forex".into(),
2926 lot_base_units: 100,
2927 lot_step_units: 1,
2928 lot_min_steps: 1,
2929 lot_max_steps: 0,
2930 }
2931 }
2932
2933 fn fixed_lot_config() -> BacktestConfig {
2934 BacktestConfig {
2935 sizing: Some(SizingPolicy::FixedLot { lots: 1.0 }),
2936 symbol_specs: ["EURUSD", "XAUUSD"]
2937 .into_iter()
2938 .map(|symbol| (symbol.to_owned(), test_symbol_spec(symbol)))
2939 .collect(),
2940 ..BacktestConfig::default()
2941 }
2942 }
2943
2944 fn identity_currency_plan(symbol: &str) -> RunCurrencyPlan {
2945 RunCurrencyPlan::new(
2946 "USD",
2947 [symbol.to_owned()].into_iter().collect(),
2948 Default::default(),
2949 [(symbol.to_owned(), "USD".to_owned())]
2950 .into_iter()
2951 .collect(),
2952 [(
2953 "USD".to_owned(),
2954 ConversionRoute::Identity {
2955 currency: "USD".to_owned(),
2956 },
2957 )]
2958 .into_iter()
2959 .collect(),
2960 Vec::new(),
2961 )
2962 .unwrap()
2963 }
2964
2965 struct ScriptedBatchFeed {
2966 batches: VecDeque<Result<Option<TimestampBatch>, &'static str>>,
2967 }
2968
2969 impl FallibleBatchFeed for ScriptedBatchFeed {
2970 type Error = &'static str;
2971
2972 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
2973 self.batches.pop_front().unwrap_or(Ok(None))
2974 }
2975 }
2976
2977 struct CountingBatchFeed {
2978 batches: VecDeque<TimestampBatch>,
2979 polls: std::rc::Rc<std::cell::Cell<usize>>,
2980 }
2981
2982 impl FallibleBatchFeed for CountingBatchFeed {
2983 type Error = Infallible;
2984
2985 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
2986 self.polls.set(self.polls.get() + 1);
2987 Ok(self.batches.pop_front())
2988 }
2989 }
2990
2991 fn primary_batch(event: MarketEvent) -> TimestampBatch {
2992 TimestampBatch {
2993 ts: event.ts(),
2994 events: vec![FeedEvent::new(
2995 event,
2996 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
2997 )],
2998 }
2999 }
3000
3001 fn market_entry(timestamp: NaiveDateTime, symbol: &str, order_type: OrderType) -> RawSignal {
3002 RawSignal::Entry {
3003 ts: timestamp,
3004 symbol: symbol.into(),
3005 side: Side::Buy,
3006 order_type,
3007 price: (order_type == OrderType::Limit).then_some(1.0),
3008 risk_multiplier: 1.0,
3009 stoploss: None,
3010 targets: Vec::new(),
3011 group: None,
3012 trade_id: Some(format!("{symbol}-blocker")),
3013 }
3014 }
3015
3016 #[test]
3017 fn future_streaming_matches_materialized_and_stops_without_draining() {
3018 let events = vec![
3019 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
3020 tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
3021 tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
3022 ];
3023 let signals = vec![
3024 market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market),
3025 RawSignal::CloseAll { ts: ts(10, 0, 1) },
3026 ];
3027 let config = BacktestConfig {
3028 close_on_finish: false,
3029 ..fixed_lot_config()
3030 };
3031 let mut materialized_feed = VecFeed::new(events.clone());
3032 let materialized = BacktestRunner::new_future(
3033 config.clone(),
3034 FutureQuoteConfig {
3035 mtm_output: MtmOutputPolicy::Full,
3036 ..FutureQuoteConfig::default()
3037 },
3038 )
3039 .run_raw_signals_future(&mut materialized_feed, signals.clone(), None);
3040
3041 let mut stream = ScriptedBatchFeed {
3042 batches: VecDeque::from([
3043 Ok(Some(primary_batch(events[0].clone()))),
3044 Ok(Some(primary_batch(events[1].clone()))),
3045 Ok(Some(primary_batch(events[2].clone()))),
3046 Err("must not drain"),
3047 ]),
3048 };
3049 let mut progress = Vec::new();
3050 let streamed = BacktestRunner::new_future(
3051 config,
3052 FutureQuoteConfig {
3053 mtm_output: MtmOutputPolicy::Full,
3054 ..FutureQuoteConfig::default()
3055 },
3056 )
3057 .run_raw_signals_future_streaming_controlled(
3058 &mut stream,
3059 Some(ts(10, 0, 2)),
3060 signals,
3061 None,
3062 || false,
3063 |update| progress.push(update),
3064 )
3065 .unwrap();
3066
3067 assert_eq!(
3068 serde_json::to_value(&streamed).unwrap(),
3069 serde_json::to_value(&materialized).unwrap()
3070 );
3071 assert_eq!(
3072 stream.batches.len(),
3073 2,
3074 "quiescence must leave the tail unread"
3075 );
3076 assert_eq!(progress.first().unwrap().total_events, 0);
3077 assert_eq!(progress.last().unwrap().processed_events, 2);
3078 assert_eq!(progress.last().unwrap().total_events, 2);
3079 assert_eq!(
3080 streamed.mtm_equity_curve.last().unwrap().ts,
3081 ts(10, 0, 1),
3082 "terminal observation must use the last processed primary timestamp"
3083 );
3084 assert_eq!(
3085 streamed
3086 .mtm_equity_curve
3087 .last()
3088 .unwrap()
3089 .observation_kind
3090 .as_deref(),
3091 Some(EquityObservationKind::QuiescentTermination.as_str())
3092 );
3093 assert_eq!(
3094 streamed
3095 .execution_metadata
3096 .as_ref()
3097 .unwrap()
3098 .tags
3099 .get("termination_reason")
3100 .map(String::as_str),
3101 Some("quiescent")
3102 );
3103 }
3104
3105 #[test]
3106 fn exact_time_close_waits_for_later_symbol_pending_fill() {
3107 let open_ts = ts(10, 0, 0);
3108 let execution_ts = ts(10, 0, 1);
3109 let events = vec![
3110 FeedEvent::new(
3111 tick("XAUUSD", 101.0, 101.0, open_ts),
3112 EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
3113 ),
3114 FeedEvent::new(
3115 tick("EURUSD", 1.1, 1.1, execution_ts),
3116 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
3117 ),
3118 FeedEvent::new(
3119 tick("XAUUSD", 100.0, 100.0, execution_ts),
3120 EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
3121 ),
3122 ];
3123 let signals = vec![
3124 RawSignal::Entry {
3125 ts: open_ts,
3126 symbol: "XAUUSD".into(),
3127 side: Side::Buy,
3128 order_type: OrderType::Limit,
3129 price: Some(100.0),
3130 risk_multiplier: 1.0,
3131 stoploss: None,
3132 targets: Vec::new(),
3133 group: None,
3134 trade_id: Some("later-pending".into()),
3135 },
3136 RawSignal::Close {
3137 ts: execution_ts,
3138 position: PositionRef::ByTradeId {
3139 trade_id: "later-pending".into(),
3140 },
3141 },
3142 ];
3143 let mut feed = VecFeed::from_feed_events(events);
3144 let result = BacktestRunner::new_future(
3145 BacktestConfig {
3146 close_on_finish: false,
3147 ..fixed_lot_config()
3148 },
3149 FutureQuoteConfig::default(),
3150 )
3151 .run_raw_signals_future(&mut feed, signals, None);
3152
3153 assert_eq!(
3154 result
3155 .recorded_fills
3156 .iter()
3157 .map(|fill| fill.fill.purpose)
3158 .collect::<Vec<_>>(),
3159 vec![FillPurpose::LimitEntry, FillPurpose::MarketExit]
3160 );
3161 assert_eq!(result.close_events.len(), 1);
3162 assert_eq!(result.close_events[0].reason, CloseReason::Manual);
3163 assert!(result.open_position_snapshots.is_empty());
3164 assert!(result.pending_order_snapshots.is_empty());
3165 }
3166
3167 #[test]
3168 fn exact_time_close_cannot_beat_later_symbol_stoploss() {
3169 let open_ts = ts(10, 0, 0);
3170 let execution_ts = ts(10, 0, 1);
3171 let events = vec![
3172 FeedEvent::new(
3173 tick("XAUUSD", 100.0, 100.0, open_ts),
3174 EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
3175 ),
3176 FeedEvent::new(
3177 tick("EURUSD", 1.1, 1.1, execution_ts),
3178 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
3179 ),
3180 FeedEvent::new(
3181 tick("XAUUSD", 98.0, 98.0, execution_ts),
3182 EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
3183 ),
3184 ];
3185 let signals = vec![
3186 RawSignal::Entry {
3187 ts: open_ts,
3188 symbol: "XAUUSD".into(),
3189 side: Side::Buy,
3190 order_type: OrderType::Market,
3191 price: None,
3192 risk_multiplier: 1.0,
3193 stoploss: Some(99.0),
3194 targets: Vec::new(),
3195 group: None,
3196 trade_id: Some("later-stop".into()),
3197 },
3198 RawSignal::Close {
3199 ts: execution_ts,
3200 position: PositionRef::ByTradeId {
3201 trade_id: "later-stop".into(),
3202 },
3203 },
3204 ];
3205 let mut feed = VecFeed::from_feed_events(events);
3206 let result = BacktestRunner::new_future(
3207 BacktestConfig {
3208 close_on_finish: false,
3209 ..fixed_lot_config()
3210 },
3211 FutureQuoteConfig::default(),
3212 )
3213 .run_raw_signals_future(&mut feed, signals, None);
3214
3215 assert_eq!(result.close_events.len(), 1);
3216 assert_eq!(result.close_events[0].reason, CloseReason::Stoploss);
3217 assert_eq!(
3218 result.recorded_fills.last().unwrap().fill.purpose,
3219 FillPurpose::StopLoss
3220 );
3221 assert!(!result.action_dispositions.iter().any(|disposition| {
3222 disposition.action_id.starts_with("signal:00000001")
3223 && disposition.status == crate::ledger::ActionDispositionStatus::Applied
3224 }));
3225 }
3226
3227 #[test]
3228 fn exact_time_multisymbol_closes_preserve_signal_order() {
3229 let open_ts = ts(10, 0, 0);
3230 let close_ts = ts(10, 0, 1);
3231 let mut events = Vec::new();
3232 for (timestamp, row) in [(open_ts, 0), (close_ts, 1)] {
3233 events.push(FeedEvent::new(
3234 tick("EURUSD", 1.1, 1.1, timestamp),
3235 EventMetadata::new(SeriesRoles::PRIMARY, 0, row),
3236 ));
3237 events.push(FeedEvent::new(
3238 tick("XAUUSD", 100.0, 100.0, timestamp),
3239 EventMetadata::new(SeriesRoles::PRIMARY, 1, row),
3240 ));
3241 }
3242 let entry = |symbol: &str, trade_id: &str| RawSignal::Entry {
3243 ts: open_ts,
3244 symbol: symbol.into(),
3245 side: Side::Buy,
3246 order_type: OrderType::Market,
3247 price: None,
3248 risk_multiplier: 1.0,
3249 stoploss: None,
3250 targets: Vec::new(),
3251 group: None,
3252 trade_id: Some(trade_id.into()),
3253 };
3254 let close = |trade_id: &str| RawSignal::Close {
3255 ts: close_ts,
3256 position: PositionRef::ByTradeId {
3257 trade_id: trade_id.into(),
3258 },
3259 };
3260 let signals = vec![
3261 entry("XAUUSD", "close-first"),
3262 entry("EURUSD", "close-second"),
3263 close("close-first"),
3264 close("close-second"),
3265 ];
3266 let mut feed = VecFeed::from_feed_events(events);
3267 let result = BacktestRunner::new_future(
3268 BacktestConfig {
3269 close_on_finish: false,
3270 ..fixed_lot_config()
3271 },
3272 FutureQuoteConfig::default(),
3273 )
3274 .run_raw_signals_future(&mut feed, signals, None);
3275
3276 assert_eq!(
3277 result
3278 .close_events
3279 .iter()
3280 .map(|event| event.symbol.as_str())
3281 .collect::<Vec<_>>(),
3282 vec!["XAUUSD", "EURUSD"]
3283 );
3284 assert!(
3285 result
3286 .close_events
3287 .iter()
3288 .all(|event| event.reason == CloseReason::Manual)
3289 );
3290 }
3291
3292 #[test]
3293 fn future_streaming_quiescence_waits_for_all_blockers() {
3294 let run = |events: Vec<MarketEvent>, signals: Vec<RawSignal>, config: BacktestConfig| {
3295 let polls = std::rc::Rc::new(std::cell::Cell::new(0));
3296 let primary_eod = events.last().map(MarketEvent::ts);
3297 let mut feed = CountingBatchFeed {
3298 batches: events.into_iter().map(primary_batch).collect(),
3299 polls: polls.clone(),
3300 };
3301 BacktestRunner::new_future(config, FutureQuoteConfig::default())
3302 .run_raw_signals_future_streaming_controlled(
3303 &mut feed,
3304 primary_eod,
3305 signals,
3306 None,
3307 || false,
3308 |_| {},
3309 )
3310 .unwrap();
3311 polls.get()
3312 };
3313 let eur_events = vec![
3314 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
3315 tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
3316 tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
3317 ];
3318
3319 let immediately_quiescent = run(
3320 eur_events.clone(),
3321 vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
3322 BacktestConfig::default(),
3323 );
3324 assert_eq!(immediately_quiescent, 1);
3325
3326 let scheduled = run(
3327 eur_events.clone(),
3328 vec![RawSignal::CloseAll { ts: ts(10, 0, 2) }],
3329 BacktestConfig::default(),
3330 );
3331 assert_eq!(scheduled, 3, "scheduled signals must block termination");
3332
3333 let mut two_symbol_config = BacktestConfig {
3334 close_on_finish: false,
3335 ..fixed_lot_config()
3336 };
3337 two_symbol_config
3338 .symbol_specs
3339 .insert("GBPUSD".into(), test_symbol_spec("GBPUSD"));
3340 let queued = run(
3341 vec![
3342 tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
3343 tick("GBPUSD", 1.2500, 1.2502, ts(10, 0, 1)),
3344 tick("GBPUSD", 1.2501, 1.2503, ts(10, 0, 2)),
3345 ],
3346 vec![
3347 market_entry(ts(10, 0, 0), "GBPUSD", OrderType::Market),
3348 RawSignal::CloseAll { ts: ts(10, 0, 1) },
3349 ],
3350 two_symbol_config,
3351 );
3352 assert_eq!(
3353 queued, 2,
3354 "queued actions must wait for an eligible symbol quote"
3355 );
3356
3357 let open = run(
3358 eur_events.clone(),
3359 vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market)],
3360 BacktestConfig {
3361 close_on_finish: false,
3362 ..fixed_lot_config()
3363 },
3364 );
3365 assert_eq!(
3366 open, 4,
3367 "open positions must consume the stream through EOD"
3368 );
3369
3370 let pending = run(
3371 eur_events,
3372 vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Limit)],
3373 fixed_lot_config(),
3374 );
3375 assert_eq!(
3376 pending, 4,
3377 "pending orders must consume the stream through EOD"
3378 );
3379 }
3380
3381 #[test]
3382 fn future_mtm_output_policies_bound_curve_and_validate_before_feed_use() {
3383 assert_eq!(
3384 FutureQuoteConfig::default().mtm_output,
3385 MtmOutputPolicy::Bounded { max_points: 4_096 }
3386 );
3387 let events: Vec<_> = (0..12)
3388 .map(|second| tick("EURUSD", 100.0, 100.0, ts(10, 0, second)))
3389 .collect();
3390
3391 let pending = RawSignal::Entry {
3392 ts: ts(10, 0, 0),
3393 symbol: "EURUSD".into(),
3394 side: Side::Buy,
3395 order_type: OrderType::Limit,
3396 price: Some(90.0),
3397 risk_multiplier: 1.0,
3398 stoploss: None,
3399 targets: Vec::new(),
3400 group: None,
3401 trade_id: Some("mtm-policy-blocker".into()),
3402 };
3403 let run = |policy| {
3404 let mut feed = VecFeed::new(events.clone());
3405 BacktestRunner::new_future(
3406 fixed_lot_config(),
3407 FutureQuoteConfig {
3408 mtm_output: policy,
3409 ..FutureQuoteConfig::default()
3410 },
3411 )
3412 .run_raw_signals_future(&mut feed, vec![pending.clone()], None)
3413 };
3414
3415 let none = run(MtmOutputPolicy::None);
3416 assert!(none.mtm_equity_curve.is_empty());
3417 assert_eq!(none.mtm_output_summary.observed_points, 13);
3418 assert_eq!(none.mtm_output_summary.omitted_points, 13);
3419
3420 let bounded = run(MtmOutputPolicy::Bounded { max_points: 8 });
3421 assert_eq!(bounded.mtm_equity_curve.len(), 8);
3422 assert_eq!(bounded.mtm_output_summary.observed_points, 13);
3423 assert_eq!(bounded.mtm_output_summary.retained_points, 8);
3424 assert_eq!(bounded.mtm_output_summary.omitted_points, 5);
3425
3426 let full = run(MtmOutputPolicy::Full);
3427 assert_eq!(full.mtm_equity_curve.len(), 13);
3428 assert_eq!(full.mtm_output_summary.observed_points, 13);
3429 assert_eq!(full.mtm_output_summary.omitted_points, 0);
3430 assert_eq!(
3431 full.mtm_equity_curve
3432 .iter()
3433 .filter(|point| {
3434 point.observation_kind.as_deref()
3435 == Some(EquityObservationKind::PostOutput.as_str())
3436 })
3437 .count(),
3438 0
3439 );
3440
3441 let mut invalid_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
3442 let rejected = BacktestRunner::new_future(
3443 BacktestConfig::default(),
3444 FutureQuoteConfig {
3445 mtm_output: MtmOutputPolicy::Bounded { max_points: 7 },
3446 ..FutureQuoteConfig::default()
3447 },
3448 )
3449 .run_raw_signals_future(&mut invalid_feed, Vec::new(), None);
3450 assert_eq!(invalid_feed.remaining(), 1);
3451 assert!(rejected.action_dispositions.iter().any(|disposition| {
3452 disposition.action_id == "configuration"
3453 && disposition
3454 .reason
3455 .as_deref()
3456 .is_some_and(|reason| reason.contains("MTM max_points"))
3457 }));
3458 }
3459
3460 #[test]
3461 fn future_mtm_records_changed_post_output_observation_kind() {
3462 let mut feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
3463 let signal = RawSignal::Entry {
3464 ts: ts(10, 0, 0),
3465 symbol: "EURUSD".into(),
3466 side: Side::Buy,
3467 order_type: OrderType::Market,
3468 price: None,
3469 risk_multiplier: 1.0,
3470 stoploss: None,
3471 targets: Vec::new(),
3472 group: None,
3473 trade_id: Some("mtm-kind".into()),
3474 };
3475 let result = BacktestRunner::new_future(
3476 BacktestConfig {
3477 close_on_finish: false,
3478 ..fixed_lot_config()
3479 },
3480 FutureQuoteConfig {
3481 mtm_output: MtmOutputPolicy::Full,
3482 ..FutureQuoteConfig::default()
3483 },
3484 )
3485 .run_raw_signals_future(&mut feed, vec![signal], None);
3486
3487 let kinds: Vec<_> = result
3488 .mtm_equity_curve
3489 .iter()
3490 .filter_map(|point| point.observation_kind.as_deref())
3491 .collect();
3492 assert_eq!(
3493 kinds,
3494 vec![
3495 EquityObservationKind::PreSettlement.as_str(),
3496 EquityObservationKind::PostOutput.as_str(),
3497 EquityObservationKind::EndOfData.as_str(),
3498 ]
3499 );
3500 assert_eq!(
3501 result
3502 .execution_metadata
3503 .as_ref()
3504 .unwrap()
3505 .tags
3506 .get("termination_reason")
3507 .map(String::as_str),
3508 Some("end_of_data")
3509 );
3510 }
3511
3512 #[test]
3513 fn future_fallible_batch_feed_propagates_source_error() {
3514 let batch = TimestampBatch {
3515 ts: ts(10, 0, 0),
3516 events: vec![FeedEvent::new(
3517 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
3518 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
3519 )],
3520 };
3521 let mut feed = ScriptedBatchFeed {
3522 batches: VecDeque::from([Ok(Some(batch)), Err("feed failed")]),
3523 };
3524 let result =
3525 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
3526 .run_raw_signals_future_fallible(&mut feed, Vec::new(), None);
3527
3528 assert!(matches!(result, Err("feed failed")));
3529 }
3530
3531 struct BuyOnceStrategy {
3535 entered: bool,
3536 }
3537
3538 impl BuyOnceStrategy {
3539 fn new() -> Self {
3540 Self { entered: false }
3541 }
3542 }
3543
3544 impl Strategy for BuyOnceStrategy {
3545 fn on_event(&mut self, event: &MarketEvent) -> Vec<Action> {
3546 if self.entered {
3547 return vec![];
3548 }
3549 if let MarketEvent::Tick { symbol, ask, .. } = event {
3550 self.entered = true;
3551 vec![Action::Open {
3552 symbol: symbol.clone(),
3553 side: Side::Buy,
3554 order_type: OrderType::Market,
3555 price: Some(*ask),
3556 size: 1.0,
3557 stoploss: Some(*ask - 0.0050),
3558 targets: vec![TargetSpec {
3559 price: *ask + 0.0050,
3560 close_ratio: 1.0,
3561 }],
3562 rules: vec![],
3563 group: None,
3564 trade_id: None,
3565 }]
3566 } else {
3567 vec![]
3568 }
3569 }
3570
3571 fn on_finished(&mut self) -> Vec<Action> {
3572 vec![]
3575 }
3576 }
3577
3578 #[test]
3581 fn strategy_backtest_tp_hit() {
3582 let events = vec![
3583 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3584 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3585 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
3586 tick("EURUSD", 1.0890, 1.0892, ts(10, 0, 3)),
3587 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 4)),
3589 ];
3590 let mut feed = VecFeed::new(events);
3591 let mut strategy = BuyOnceStrategy::new();
3592
3593 let config = BacktestConfig {
3594 initial_balance: 10_000.0,
3595 close_on_finish: true,
3596 ..Default::default()
3597 };
3598 let runner = BacktestRunner::new(config);
3599 let result = runner.run_strategy(&mut feed, &mut strategy);
3600
3601 assert_eq!(result.total_trades, 1);
3602 assert_eq!(result.winning_trades, 1);
3603 assert!(result.total_pnl > 0.0);
3604 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
3605 }
3606
3607 #[test]
3608 fn strategy_backtest_sl_hit() {
3609 let events = vec![
3610 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3611 tick("EURUSD", 1.0830, 1.0832, ts(10, 0, 1)),
3612 tick("EURUSD", 1.0799, 1.0801, ts(10, 0, 2)),
3614 ];
3615 let mut feed = VecFeed::new(events);
3616 let mut strategy = BuyOnceStrategy::new();
3617
3618 let runner = BacktestRunner::with_defaults();
3619 let result = runner.run_strategy(&mut feed, &mut strategy);
3620
3621 assert_eq!(result.total_trades, 1);
3622 assert_eq!(result.losing_trades, 1);
3623 assert!(result.total_pnl < 0.0);
3624 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
3625 }
3626
3627 #[test]
3628 fn strategy_close_on_finish() {
3629 let events = vec![
3631 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3632 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3633 tick("EURUSD", 1.0852, 1.0854, ts(10, 0, 2)),
3634 ];
3635 let mut feed = VecFeed::new(events);
3636 let mut strategy = BuyOnceStrategy::new();
3637
3638 let config = BacktestConfig {
3639 initial_balance: 10_000.0,
3640 close_on_finish: true,
3641 ..Default::default()
3642 };
3643 let runner = BacktestRunner::new(config);
3644 let result = runner.run_strategy(&mut feed, &mut strategy);
3645
3646 assert_eq!(result.total_trades, 1);
3647 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
3648 }
3649
3650 #[test]
3651 fn strategy_no_close_on_finish() {
3652 let events = vec![
3653 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3654 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3655 ];
3656 let mut feed = VecFeed::new(events);
3657 let mut strategy = BuyOnceStrategy::new();
3658
3659 let config = BacktestConfig {
3660 initial_balance: 10_000.0,
3661 close_on_finish: false,
3662 ..Default::default()
3663 };
3664 let runner = BacktestRunner::new(config);
3665 let result = runner.run_strategy(&mut feed, &mut strategy);
3666
3667 assert_eq!(result.total_trades, 0);
3669 }
3670
3671 #[test]
3674 fn legacy_unprofiled_targets_default_to_equal_weights() {
3675 let events = vec![
3676 tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
3677 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
3678 tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
3679 ];
3680 let mut feed = VecFeed::new(events);
3681 let signals = vec![RawSignal::Entry {
3682 ts: ts(10, 0, 0),
3683 symbol: "EURUSD".into(),
3684 side: Side::Buy,
3685 order_type: OrderType::Market,
3686 price: Some(1.0000),
3687 risk_multiplier: 1.0,
3688 stoploss: None,
3689 targets: vec![1.1000, 1.2000],
3690 group: None,
3691 trade_id: Some("equal-targets".into()),
3692 }];
3693
3694 let result = BacktestRunner::new(BacktestConfig {
3695 close_on_finish: false,
3696 ..fixed_lot_config()
3697 })
3698 .run_raw_signals(&mut feed, signals, None);
3699
3700 assert_eq!(result.trade_log.len(), 2);
3701 assert!(
3702 result
3703 .trade_log
3704 .iter()
3705 .all(|trade| (trade.size - 0.5).abs() < f64::EPSILON)
3706 );
3707 assert!(
3708 result
3709 .trade_log
3710 .iter()
3711 .all(|trade| trade.close_reason == CloseReason::Target)
3712 );
3713 }
3714
3715 #[test]
3716 fn legacy_atomic_target_modification_retains_profile_ratio() {
3717 let events = vec![
3718 tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
3719 tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
3720 tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
3721 tick("EURUSD", 1.3000, 1.3000, ts(10, 0, 3)),
3722 ];
3723 let mut feed = VecFeed::new(events);
3724 let position = PositionRef::ByTradeId {
3725 trade_id: "modified-target".into(),
3726 };
3727 let signals = vec![
3728 RawSignal::Entry {
3729 ts: ts(10, 0, 0),
3730 symbol: "EURUSD".into(),
3731 side: Side::Buy,
3732 order_type: OrderType::Market,
3733 price: Some(1.0000),
3734 risk_multiplier: 1.0,
3735 stoploss: None,
3736 targets: vec![1.1000, 1.3000],
3737 group: None,
3738 trade_id: Some("modified-target".into()),
3739 },
3740 RawSignal::ModifyTarget {
3741 ts: ts(10, 0, 0),
3742 position,
3743 old_price: 1.1000,
3744 new_price: 1.2000,
3745 },
3746 ];
3747 let profile = ManagementProfile {
3748 name: "non-default-ratios".into(),
3749 target_selection: None,
3750 use_targets: vec![1, 2],
3751 close_ratios: vec![0.25, 0.75],
3752 stoploss_mode: StoplossMode::FromSignal,
3753 rules: vec![],
3754 group_override: None,
3755 let_remainder_run: false,
3756 };
3757
3758 let result = BacktestRunner::new(BacktestConfig {
3759 close_on_finish: false,
3760 ..fixed_lot_config()
3761 })
3762 .run_raw_signals(&mut feed, signals, Some(&profile));
3763
3764 assert_eq!(result.trade_log.len(), 2);
3765 assert!((result.trade_log[0].exit_price - 1.2000).abs() < f64::EPSILON);
3766 assert!((result.trade_log[0].size - 0.25).abs() < f64::EPSILON);
3767 assert!((result.trade_log[1].exit_price - 1.3000).abs() < f64::EPSILON);
3768 assert!((result.trade_log[1].size - 0.75).abs() < f64::EPSILON);
3769 }
3770
3771 #[test]
3772 fn run_raw_signals_entry_only() {
3773 let events = vec![
3774 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3775 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3776 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
3777 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
3778 ];
3779 let mut feed = VecFeed::new(events);
3780
3781 let raw_signals = vec![RawSignal::Entry {
3782 ts: ts(10, 0, 0),
3783 symbol: "EURUSD".into(),
3784 side: Side::Buy,
3785 order_type: OrderType::Market,
3786 price: Some(1.0850),
3787 risk_multiplier: 1.0,
3788 stoploss: Some(1.0800),
3789 targets: vec![1.0900],
3790 group: None,
3791 trade_id: None,
3792 }];
3793
3794 let runner = BacktestRunner::new(fixed_lot_config());
3795 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
3796
3797 assert_eq!(result.total_trades, 1);
3798 assert_eq!(result.winning_trades, 1);
3799 }
3800
3801 #[test]
3802 fn run_raw_signals_open_then_close() {
3803 let events = vec![
3804 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3805 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3806 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
3807 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
3808 ];
3809 let mut feed = VecFeed::new(events);
3810
3811 let raw_signals = vec![
3812 RawSignal::Entry {
3813 ts: ts(10, 0, 0),
3814 symbol: "EURUSD".into(),
3815 side: Side::Buy,
3816 order_type: OrderType::Market,
3817 price: Some(1.0850),
3818 risk_multiplier: 1.0,
3819 stoploss: None,
3820 targets: vec![],
3821 group: None,
3822 trade_id: Some("t1".into()),
3823 },
3824 RawSignal::Close {
3825 ts: ts(10, 0, 2),
3826 position: PositionRef::ByTradeId {
3827 trade_id: "t1".into(),
3828 },
3829 },
3830 ];
3831
3832 let config = BacktestConfig {
3833 initial_balance: 10_000.0,
3834 close_on_finish: false,
3835 ..fixed_lot_config()
3836 };
3837 let runner = BacktestRunner::new(config);
3838 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
3839
3840 assert_eq!(result.total_trades, 1);
3841 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
3842 }
3843
3844 #[test]
3845 fn run_raw_signals_open_then_modify_sl() {
3846 let events = vec![
3848 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3849 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
3850 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 2)),
3852 tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
3854 ];
3855 let mut feed = VecFeed::new(events);
3856
3857 let raw_signals = vec![
3858 RawSignal::Entry {
3859 ts: ts(10, 0, 0),
3860 symbol: "EURUSD".into(),
3861 side: Side::Buy,
3862 order_type: OrderType::Market,
3863 price: Some(1.0850),
3864 risk_multiplier: 1.0,
3865 stoploss: Some(1.0800),
3866 targets: vec![],
3867 group: None,
3868 trade_id: Some("t1".into()),
3869 },
3870 RawSignal::ModifyStoploss {
3871 ts: ts(10, 0, 2),
3872 position: PositionRef::ByTradeId {
3873 trade_id: "t1".into(),
3874 },
3875 price: 1.0840,
3876 },
3877 ];
3878
3879 let config = BacktestConfig {
3880 initial_balance: 10_000.0,
3881 close_on_finish: true,
3882 ..fixed_lot_config()
3883 };
3884 let runner = BacktestRunner::new(config);
3885 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
3886
3887 assert_eq!(result.total_trades, 1);
3888 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
3889 }
3890
3891 #[test]
3892 fn run_raw_signals_open_then_partial_close() {
3893 let events = vec![
3894 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3895 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
3896 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
3897 tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 3)),
3898 ];
3899 let mut feed = VecFeed::new(events);
3900
3901 let raw_signals = vec![
3902 RawSignal::Entry {
3903 ts: ts(10, 0, 0),
3904 symbol: "EURUSD".into(),
3905 side: Side::Buy,
3906 order_type: OrderType::Market,
3907 price: Some(1.0850),
3908 risk_multiplier: 1.0,
3909 stoploss: None,
3910 targets: vec![],
3911 group: None,
3912 trade_id: Some("t1".into()),
3913 },
3914 RawSignal::ClosePartial {
3915 ts: ts(10, 0, 1),
3916 position: PositionRef::ByTradeId {
3917 trade_id: "t1".into(),
3918 },
3919 ratio: 0.5,
3920 },
3921 ];
3922
3923 let config = BacktestConfig {
3924 initial_balance: 10_000.0,
3925 close_on_finish: true,
3926 ..fixed_lot_config()
3927 };
3928 let runner = BacktestRunner::new(config);
3929 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
3930
3931 assert!(result.total_trades >= 1);
3933 }
3934
3935 #[test]
3936 fn run_raw_signals_group_workflow() {
3937 let events = vec![
3938 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3939 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3940 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
3941 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
3942 tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 4)),
3943 ];
3944 let mut feed = VecFeed::new(events);
3945
3946 let raw_signals = vec![
3947 RawSignal::Entry {
3949 ts: ts(10, 0, 0),
3950 symbol: "EURUSD".into(),
3951 side: Side::Buy,
3952 order_type: OrderType::Market,
3953 price: Some(1.0850),
3954 risk_multiplier: 1.0,
3955 stoploss: None,
3956 targets: vec![],
3957 group: Some("grp1".into()),
3958 trade_id: Some("t1".into()),
3959 },
3960 RawSignal::Entry {
3961 ts: ts(10, 0, 1),
3962 symbol: "EURUSD".into(),
3963 side: Side::Buy,
3964 order_type: OrderType::Market,
3965 price: Some(1.0857),
3966 risk_multiplier: 1.0,
3967 stoploss: None,
3968 targets: vec![],
3969 group: Some("grp1".into()),
3970 trade_id: Some("t2".into()),
3971 },
3972 RawSignal::CloseAllInGroup {
3974 ts: ts(10, 0, 3),
3975 group_id: "grp1".into(),
3976 },
3977 ];
3978
3979 let config = BacktestConfig {
3980 initial_balance: 10_000.0,
3981 close_on_finish: false,
3982 ..fixed_lot_config()
3983 };
3984 let runner = BacktestRunner::new(config);
3985 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
3986
3987 assert_eq!(result.total_trades, 2);
3988 }
3989
3990 #[test]
3991 fn run_raw_signals_close_all_of_symbol() {
3992 let events = vec![
3993 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
3994 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
3995 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
3996 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
3997 ];
3998 let mut feed = VecFeed::new(events);
3999
4000 let raw_signals = vec![
4001 RawSignal::Entry {
4002 ts: ts(10, 0, 0),
4003 symbol: "EURUSD".into(),
4004 side: Side::Buy,
4005 order_type: OrderType::Market,
4006 price: Some(1.0850),
4007 risk_multiplier: 1.0,
4008 stoploss: None,
4009 targets: vec![],
4010 group: None,
4011 trade_id: Some("t1".into()),
4012 },
4013 RawSignal::Entry {
4014 ts: ts(10, 0, 0),
4015 symbol: "EURUSD".into(),
4016 side: Side::Buy,
4017 order_type: OrderType::Market,
4018 price: Some(1.0850),
4019 risk_multiplier: 0.5,
4020 stoploss: None,
4021 targets: vec![],
4022 group: None,
4023 trade_id: Some("t2".into()),
4024 },
4025 RawSignal::CloseAllOf {
4026 ts: ts(10, 0, 2),
4027 symbol: "EURUSD".into(),
4028 },
4029 ];
4030
4031 let config = BacktestConfig {
4032 initial_balance: 10_000.0,
4033 close_on_finish: false,
4034 ..fixed_lot_config()
4035 };
4036 let runner = BacktestRunner::new(config);
4037 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4038
4039 assert_eq!(result.total_trades, 2);
4040 }
4041
4042 #[test]
4043 fn run_raw_signals_with_profile() {
4044 let events = vec![
4045 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4046 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4047 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
4048 ];
4049 let mut feed = VecFeed::new(events);
4050
4051 let profile = ManagementProfile {
4052 name: "test".into(),
4053 target_selection: None,
4054 use_targets: vec![1],
4055 close_ratios: vec![1.0],
4056 stoploss_mode: StoplossMode::FromSignal,
4057 rules: vec![],
4058 group_override: None,
4059 let_remainder_run: false,
4060 };
4061
4062 let raw_signals = vec![RawSignal::Entry {
4063 ts: ts(10, 0, 0),
4064 symbol: "EURUSD".into(),
4065 side: Side::Buy,
4066 order_type: OrderType::Market,
4067 price: Some(1.0850),
4068 risk_multiplier: 1.0,
4069 stoploss: Some(1.0800),
4070 targets: vec![1.0900],
4071 group: None,
4072 trade_id: Some("t1".into()),
4073 }];
4074
4075 let runner = BacktestRunner::new(fixed_lot_config());
4076 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
4077
4078 assert_eq!(result.total_trades, 1);
4079 assert_eq!(result.winning_trades, 1);
4080 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
4081 }
4082
4083 #[test]
4084 fn run_raw_signals_with_profile_preserves_trade_id() {
4085 let events = vec![
4088 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4089 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
4090 ];
4091 let mut feed = VecFeed::new(events);
4092
4093 let profile = ManagementProfile {
4094 name: "test".into(),
4095 target_selection: None,
4096 use_targets: vec![1],
4097 close_ratios: vec![1.0],
4098 stoploss_mode: StoplossMode::FromSignal,
4099 rules: vec![],
4100 group_override: None,
4101 let_remainder_run: false,
4102 };
4103
4104 let raw_signals = vec![
4105 RawSignal::Entry {
4106 ts: ts(10, 0, 0),
4107 symbol: "EURUSD".into(),
4108 side: Side::Buy,
4109 order_type: OrderType::Market,
4110 price: Some(1.0850),
4111 risk_multiplier: 1.0,
4112 stoploss: Some(1.0800),
4113 targets: vec![1.0900],
4114 group: None,
4115 trade_id: Some("msg-100".into()),
4116 },
4117 RawSignal::Close {
4118 ts: ts(10, 0, 1),
4119 position: PositionRef::ByTradeId {
4120 trade_id: "msg-100".into(),
4121 },
4122 },
4123 ];
4124
4125 let runner = BacktestRunner::new(fixed_lot_config());
4126 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
4127
4128 assert_eq!(result.total_trades, 1);
4129 assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
4130 }
4131
4132 #[test]
4133 fn run_raw_signals_no_profile() {
4134 let events = vec![
4136 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4137 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4138 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
4139 ];
4140 let mut feed = VecFeed::new(events);
4141
4142 let raw_signals = vec![RawSignal::Entry {
4143 ts: ts(10, 0, 0),
4144 symbol: "EURUSD".into(),
4145 side: Side::Buy,
4146 order_type: OrderType::Market,
4147 price: Some(1.0850),
4148 risk_multiplier: 1.0,
4149 stoploss: None,
4150 targets: vec![],
4151 group: None,
4152 trade_id: None,
4153 }];
4154
4155 let config = BacktestConfig {
4156 initial_balance: 10_000.0,
4157 close_on_finish: true,
4158 ..fixed_lot_config()
4159 };
4160 let runner = BacktestRunner::new(config);
4161 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4162
4163 assert_eq!(result.total_trades, 1);
4164 }
4165
4166 #[test]
4167 fn run_raw_signals_last_on_symbol_resolution() {
4168 let events = vec![
4170 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4171 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4172 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
4173 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
4174 ];
4175 let mut feed = VecFeed::new(events);
4176
4177 let raw_signals = vec![
4178 RawSignal::Entry {
4179 ts: ts(10, 0, 0),
4180 symbol: "EURUSD".into(),
4181 side: Side::Buy,
4182 order_type: OrderType::Market,
4183 price: Some(1.0850),
4184 risk_multiplier: 1.0,
4185 stoploss: None,
4186 targets: vec![],
4187 group: None,
4188 trade_id: Some("t1".into()),
4189 },
4190 RawSignal::Entry {
4191 ts: ts(10, 0, 1),
4192 symbol: "EURUSD".into(),
4193 side: Side::Buy,
4194 order_type: OrderType::Market,
4195 price: Some(1.0857),
4196 risk_multiplier: 1.0,
4197 stoploss: None,
4198 targets: vec![],
4199 group: None,
4200 trade_id: Some("t2".into()),
4201 },
4202 RawSignal::Close {
4204 ts: ts(10, 0, 2),
4205 position: PositionRef::ByTradeId {
4206 trade_id: "t2".into(),
4207 },
4208 },
4209 ];
4210
4211 let config = BacktestConfig {
4212 initial_balance: 10_000.0,
4213 close_on_finish: true,
4214 ..fixed_lot_config()
4215 };
4216 let runner = BacktestRunner::new(config);
4217 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4218
4219 assert_eq!(result.total_trades, 2);
4221 }
4222
4223 #[test]
4224 fn run_raw_signals_unresolved_ref_skipped() {
4225 let events = vec![
4227 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4228 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4229 ];
4230 let mut feed = VecFeed::new(events);
4231
4232 let raw_signals = vec![RawSignal::Close {
4233 ts: ts(10, 0, 0),
4234 position: PositionRef::ByTradeId {
4235 trade_id: "nonexistent".into(),
4236 },
4237 }];
4238
4239 let config = BacktestConfig {
4240 initial_balance: 10_000.0,
4241 close_on_finish: false,
4242 ..Default::default()
4243 };
4244 let runner = BacktestRunner::new(config);
4245 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4246
4247 assert_eq!(result.total_trades, 0);
4249 }
4250
4251 #[test]
4254 fn signal_replay_basic() {
4255 let events = vec![
4256 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4257 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4258 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
4259 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
4261 ];
4262 let mut feed = VecFeed::new(events);
4263
4264 let raw_signals = vec![RawSignal::Entry {
4265 ts: ts(10, 0, 0),
4266 symbol: "EURUSD".into(),
4267 side: Side::Buy,
4268 order_type: OrderType::Market,
4269 price: Some(1.0850),
4270 risk_multiplier: 1.0,
4271 stoploss: Some(1.0800),
4272 targets: vec![1.0900],
4273 group: None,
4274 trade_id: None,
4275 }];
4276
4277 let runner = BacktestRunner::new(fixed_lot_config());
4278 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4279
4280 assert_eq!(result.total_trades, 1);
4281 assert_eq!(result.winning_trades, 1);
4282 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
4283 }
4284
4285 #[test]
4286 fn signal_replay_multiple_signals() {
4287 let events = vec![
4288 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4289 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4290 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
4292 tick("EURUSD", 1.0910, 1.0912, ts(10, 0, 3)),
4293 tick("EURUSD", 1.0920, 1.0922, ts(10, 0, 4)),
4294 ];
4295 let mut feed = VecFeed::new(events);
4296
4297 let raw_signals = vec![
4298 RawSignal::Entry {
4299 ts: ts(10, 0, 0),
4300 symbol: "EURUSD".into(),
4301 side: Side::Buy,
4302 order_type: OrderType::Market,
4303 price: Some(1.0850),
4304 risk_multiplier: 1.0,
4305 stoploss: Some(1.0800),
4306 targets: vec![1.0900],
4307 group: None,
4308 trade_id: Some("t1".into()),
4309 },
4310 RawSignal::Entry {
4311 ts: ts(10, 0, 1),
4312 symbol: "EURUSD".into(),
4313 side: Side::Buy,
4314 order_type: OrderType::Market,
4315 price: Some(1.0857),
4316 risk_multiplier: 1.0,
4317 stoploss: None,
4318 targets: vec![],
4319 group: None,
4320 trade_id: Some("t2".into()),
4321 },
4322 ];
4323
4324 let config = BacktestConfig {
4325 initial_balance: 10_000.0,
4326 close_on_finish: true,
4327 ..fixed_lot_config()
4328 };
4329 let runner = BacktestRunner::new(config);
4330 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4331
4332 assert!(result.total_trades >= 2);
4334 }
4335
4336 #[test]
4337 fn signal_replay_signal_before_data_filtered() {
4338 let events = vec![
4345 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4346 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
4347 ];
4348 let mut feed = VecFeed::new(events);
4349
4350 let raw_signals = vec![RawSignal::Entry {
4351 ts: ts(9, 0, 0), symbol: "EURUSD".into(),
4353 side: Side::Buy,
4354 order_type: OrderType::Market,
4355 price: Some(1.0850),
4356 risk_multiplier: 1.0,
4357 stoploss: None,
4358 targets: vec![1.0900],
4359 group: None,
4360 trade_id: None,
4361 }];
4362
4363 let runner = BacktestRunner::new(fixed_lot_config());
4364 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4365
4366 assert_eq!(result.total_trades, 1);
4367 assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
4368 }
4369
4370 #[test]
4371 fn empty_feed_empty_result() {
4372 let mut feed = VecFeed::new(vec![]);
4373 let mut strategy = BuyOnceStrategy::new();
4374
4375 let runner = BacktestRunner::with_defaults();
4376 let result = runner.run_strategy(&mut feed, &mut strategy);
4377
4378 assert_eq!(result.total_trades, 0);
4379 assert!((result.final_balance - 10_000.0).abs() < f64::EPSILON);
4380 }
4381
4382 #[test]
4383 fn report_display_does_not_panic() {
4384 let events = vec![
4385 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4386 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
4387 ];
4388 let mut feed = VecFeed::new(events);
4389 let mut strategy = BuyOnceStrategy::new();
4390
4391 let runner = BacktestRunner::with_defaults();
4392 let result = runner.run_strategy(&mut feed, &mut strategy);
4393
4394 let _display = format!("{result}");
4395 }
4396
4397 #[test]
4398 fn run_raw_signals_with_profile_open_then_modify_sl_by_trade_id() {
4399 let events = vec![
4400 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4401 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4402 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
4403 tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
4404 ];
4405 let mut feed = VecFeed::new(events);
4406
4407 let profile = ManagementProfile {
4408 name: "test".into(),
4409 target_selection: None,
4410 use_targets: vec![1],
4411 close_ratios: vec![1.0],
4412 stoploss_mode: StoplossMode::FromSignal,
4413 rules: vec![],
4414 group_override: None,
4415 let_remainder_run: false,
4416 };
4417
4418 let raw_signals = vec![
4419 RawSignal::Entry {
4420 ts: ts(10, 0, 0),
4421 symbol: "EURUSD".into(),
4422 side: Side::Buy,
4423 order_type: OrderType::Market,
4424 price: Some(1.0850),
4425 risk_multiplier: 1.0,
4426 stoploss: Some(1.0800),
4427 targets: vec![1.0900],
4428 group: None,
4429 trade_id: Some("t1".into()),
4430 },
4431 RawSignal::ModifyStoploss {
4432 ts: ts(10, 0, 2),
4433 position: PositionRef::ByTradeId {
4434 trade_id: "t1".into(),
4435 },
4436 price: 1.0840,
4437 },
4438 ];
4439
4440 let config = BacktestConfig {
4441 initial_balance: 10_000.0,
4442 close_on_finish: false,
4443 ..fixed_lot_config()
4444 };
4445 let runner = BacktestRunner::new(config);
4446 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
4447
4448 assert_eq!(result.total_trades, 1);
4449 assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
4450 }
4451
4452 #[test]
4453 fn run_raw_signals_with_profile_open_then_close_partial_by_trade_id() {
4454 let events = vec![
4455 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4456 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4457 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
4458 ];
4459 let mut feed = VecFeed::new(events);
4460
4461 let profile = ManagementProfile {
4462 name: "test".into(),
4463 target_selection: None,
4464 use_targets: vec![1],
4465 close_ratios: vec![1.0],
4466 stoploss_mode: StoplossMode::FromSignal,
4467 rules: vec![],
4468 group_override: None,
4469 let_remainder_run: false,
4470 };
4471
4472 let raw_signals = vec![
4473 RawSignal::Entry {
4474 ts: ts(10, 0, 0),
4475 symbol: "EURUSD".into(),
4476 side: Side::Buy,
4477 order_type: OrderType::Market,
4478 price: Some(1.0850),
4479 risk_multiplier: 1.0,
4480 stoploss: Some(1.0800),
4481 targets: vec![1.0900],
4482 group: None,
4483 trade_id: Some("t1".into()),
4484 },
4485 RawSignal::ClosePartial {
4486 ts: ts(10, 0, 1),
4487 position: PositionRef::ByTradeId {
4488 trade_id: "t1".into(),
4489 },
4490 ratio: 0.5,
4491 },
4492 ];
4493
4494 let config = BacktestConfig {
4495 initial_balance: 10_000.0,
4496 close_on_finish: false,
4497 ..fixed_lot_config()
4498 };
4499 let runner = BacktestRunner::new(config);
4500 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
4501
4502 assert!(result.total_trades >= 1);
4504 }
4505
4506 #[test]
4507 fn run_raw_signals_multi_position_by_trade_id_with_profile() {
4508 let events = vec![
4512 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
4513 tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
4514 tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
4515 tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
4516 ];
4517 let mut feed = VecFeed::new(events);
4518
4519 let profile = ManagementProfile {
4520 name: "test".into(),
4521 target_selection: None,
4522 use_targets: vec![1],
4523 close_ratios: vec![1.0],
4524 stoploss_mode: StoplossMode::FromSignal,
4525 rules: vec![],
4526 group_override: Some("alpha".into()),
4527 let_remainder_run: false,
4528 };
4529
4530 let raw_signals = vec![
4531 RawSignal::Entry {
4532 ts: ts(10, 0, 0),
4533 symbol: "EURUSD".into(),
4534 side: Side::Buy,
4535 order_type: OrderType::Market,
4536 price: Some(1.0850),
4537 risk_multiplier: 1.0,
4538 stoploss: None,
4539 targets: vec![1.0910],
4540 group: None,
4541 trade_id: Some("t1".into()),
4542 },
4543 RawSignal::Entry {
4544 ts: ts(10, 0, 1),
4545 symbol: "EURUSD".into(),
4546 side: Side::Buy,
4547 order_type: OrderType::Market,
4548 price: Some(1.0857),
4549 risk_multiplier: 1.0,
4550 stoploss: None,
4551 targets: vec![1.0910],
4552 group: None,
4553 trade_id: Some("t2".into()),
4554 },
4555 RawSignal::Close {
4557 ts: ts(10, 0, 2),
4558 position: PositionRef::ByTradeId {
4559 trade_id: "t1".into(),
4560 },
4561 },
4562 ];
4563
4564 let config = BacktestConfig {
4565 initial_balance: 10_000.0,
4566 close_on_finish: true,
4567 ..fixed_lot_config()
4568 };
4569 let runner = BacktestRunner::new(config);
4570 let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
4571
4572 assert_eq!(result.total_trades, 2);
4574 for trade in &result.trade_log {
4576 assert_eq!(trade.group.as_deref(), Some("alpha"));
4577 }
4578 }
4579
4580 #[test]
4581 fn merged_feed_manual_close_uses_correct_symbol_quote() {
4582 use crate::data_feed::MarketEvent;
4587 let events = vec![
4588 MarketEvent::Tick {
4589 symbol: "XAUUSD".into(),
4590 ts: ts(10, 0, 0),
4591 bid: 5000.0,
4592 ask: 5001.0,
4593 },
4594 MarketEvent::Tick {
4595 symbol: "GBPJPY".into(),
4596 ts: ts(10, 0, 1),
4597 bid: 210.0,
4598 ask: 211.0,
4599 },
4600 MarketEvent::Tick {
4601 symbol: "XAUUSD".into(),
4602 ts: ts(10, 0, 2),
4603 bid: 5050.0,
4604 ask: 5051.0,
4605 },
4606 MarketEvent::Tick {
4608 symbol: "GBPJPY".into(),
4609 ts: ts(10, 0, 3),
4610 bid: 212.0,
4611 ask: 213.0,
4612 },
4613 ];
4614 let mut feed = VecFeed::new(events);
4615
4616 let raw_signals = vec![
4617 RawSignal::Entry {
4618 ts: ts(10, 0, 0),
4619 symbol: "XAUUSD".into(),
4620 side: Side::Buy,
4621 order_type: OrderType::Market,
4622 price: Some(5000.0),
4623 risk_multiplier: 1.0,
4624 stoploss: None,
4625 targets: vec![],
4626 group: None,
4627 trade_id: Some("xau-1".into()),
4628 },
4629 RawSignal::Close {
4631 ts: ts(10, 0, 3),
4632 position: PositionRef::ByTradeId {
4633 trade_id: "xau-1".into(),
4634 },
4635 },
4636 ];
4637
4638 let config = BacktestConfig {
4639 initial_balance: 10_000.0,
4640 close_on_finish: false,
4641 ..fixed_lot_config()
4642 };
4643 let runner = BacktestRunner::new(config);
4644 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
4645
4646 assert_eq!(result.total_trades, 1);
4647 let trade = &result.trade_log[0];
4648 assert_eq!(trade.symbol, "XAUUSD");
4649 assert!(
4651 trade.exit_price > 4000.0,
4652 "Exit price should be XAUUSD (~5050), got {}",
4653 trade.exit_price
4654 );
4655 }
4656
4657 fn long_tick_feed(count: usize) -> VecFeed {
4658 let start = ts(10, 0, 0);
4659 VecFeed::new(
4660 (0..count)
4661 .map(|index| {
4662 tick(
4663 "EURUSD",
4664 1.0848,
4665 1.0850,
4666 start + Duration::milliseconds(index as i64),
4667 )
4668 })
4669 .collect(),
4670 )
4671 }
4672
4673 #[test]
4674 fn legacy_replay_can_be_cancelled_during_event_processing() {
4675 let cancelled = std::cell::Cell::new(false);
4676 let mut feed = long_tick_feed(1_000);
4677 let outcome = BacktestRunner::with_defaults().run_raw_signals_controlled(
4678 &mut feed,
4679 Vec::new(),
4680 None,
4681 || cancelled.get(),
4682 |progress| {
4683 if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
4684 cancelled.set(true);
4685 }
4686 },
4687 );
4688
4689 assert_eq!(outcome.unwrap_err(), ReplayCancelled);
4690 assert!(
4691 feed.remaining() > 0,
4692 "cancellation must stop further replay"
4693 );
4694 }
4695
4696 #[test]
4697 fn future_quote_replay_can_be_cancelled_during_event_processing() {
4698 let cancelled = std::cell::Cell::new(false);
4699 let mut feed = long_tick_feed(1_000);
4700 let runner = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default());
4701 let pending = RawSignal::Entry {
4702 ts: ts(10, 0, 0),
4703 symbol: "EURUSD".into(),
4704 side: Side::Buy,
4705 order_type: OrderType::Limit,
4706 price: Some(1.0),
4707 risk_multiplier: 1.0,
4708 stoploss: None,
4709 targets: Vec::new(),
4710 group: None,
4711 trade_id: Some("cancellation-blocker".into()),
4712 };
4713 let outcome = runner.run_raw_signals_controlled(
4714 &mut feed,
4715 vec![pending],
4716 None,
4717 || cancelled.get(),
4718 |progress| {
4719 if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
4720 cancelled.set(true);
4721 }
4722 },
4723 );
4724
4725 assert_eq!(outcome.unwrap_err(), ReplayCancelled);
4726 }
4727
4728 #[test]
4729 fn controlled_replay_progress_is_monotonic_and_reaches_event_total() {
4730 let mut feed = long_tick_feed(600);
4731 let mut updates = Vec::new();
4732 BacktestRunner::with_defaults()
4733 .run_raw_signals_controlled(
4734 &mut feed,
4735 Vec::new(),
4736 None,
4737 || false,
4738 |progress| updates.push(progress),
4739 )
4740 .unwrap();
4741
4742 assert!(updates.len() >= 3);
4743 assert!(updates.windows(2).all(|pair| {
4744 pair[0].processed_events <= pair[1].processed_events
4745 && pair[0].processed_signals <= pair[1].processed_signals
4746 && pair[0].total_events <= pair[1].total_events
4747 && pair[0].total_signals <= pair[1].total_signals
4748 }));
4749 assert_eq!(updates.last().unwrap().processed_events, 600);
4750 assert_eq!(updates.last().unwrap().total_events, 600);
4751 }
4752
4753 #[test]
4754 fn legacy_replay_skips_invalid_crossed_and_reversed_quotes_without_nonfinite_pnl() {
4755 let events = vec![
4756 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
4757 tick("EURUSD", f64::NAN, 101.0, ts(10, 0, 1)),
4758 tick("EURUSD", 102.0, 101.0, ts(10, 0, 2)),
4759 tick("EURUSD", 90.0, 90.0, ts(9, 59, 59)),
4760 tick("EURUSD", 110.0, 110.0, ts(10, 0, 3)),
4761 ];
4762 let mut feed = VecFeed::new(events);
4763 let signals = vec![RawSignal::Entry {
4764 ts: ts(10, 0, 0),
4765 symbol: "EURUSD".into(),
4766 side: Side::Buy,
4767 order_type: OrderType::Market,
4768 price: Some(100.0),
4769 risk_multiplier: 1.0,
4770 stoploss: None,
4771 targets: vec![],
4772 group: None,
4773 trade_id: Some("safe-feed".into()),
4774 }];
4775
4776 let result =
4777 BacktestRunner::new(fixed_lot_config()).run_raw_signals(&mut feed, signals, None);
4778 assert_eq!(result.trade_log.len(), 1);
4779 assert_eq!(result.trade_log[0].exit_price, 110.0);
4780 assert_eq!(result.trade_log[0].pnl, 10.0);
4781 assert!(result.total_pnl.is_finite());
4782 assert!(result.final_balance.is_finite());
4783 }
4784
4785 #[test]
4786 fn legacy_and_future_profile_replay_share_empty_ratio_target_resolution() {
4787 let profile = ManagementProfile {
4788 name: "equal-target".into(),
4789 target_selection: None,
4790 use_targets: vec![1],
4791 close_ratios: vec![],
4792 stoploss_mode: StoplossMode::FromSignal,
4793 rules: vec![],
4794 group_override: None,
4795 let_remainder_run: false,
4796 };
4797 let signals = vec![RawSignal::Entry {
4798 ts: ts(10, 0, 0),
4799 symbol: "EURUSD".into(),
4800 side: Side::Buy,
4801 order_type: OrderType::Market,
4802 price: Some(100.0),
4803 risk_multiplier: 1.0,
4804 stoploss: None,
4805 targets: vec![101.0],
4806 group: None,
4807 trade_id: Some("profile-parity".into()),
4808 }];
4809 let events = vec![
4810 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
4811 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
4812 ];
4813
4814 let mut legacy_feed = VecFeed::new(events.clone());
4815 let legacy = BacktestRunner::new(BacktestConfig {
4816 close_on_finish: false,
4817 ..fixed_lot_config()
4818 })
4819 .run_raw_signals(&mut legacy_feed, signals.clone(), Some(&profile));
4820 let mut future_feed = VecFeed::new(events);
4821 let future = BacktestRunner::new_future(
4822 BacktestConfig {
4823 close_on_finish: false,
4824 ..fixed_lot_config()
4825 },
4826 FutureQuoteConfig::default(),
4827 )
4828 .run_raw_signals_future(&mut future_feed, signals, Some(&profile));
4829
4830 assert_eq!(legacy.trade_log.len(), 1);
4831 assert_eq!(future.trade_log.len(), 1);
4832 assert_eq!(legacy.trade_log[0].close_reason, CloseReason::Target);
4833 assert_eq!(future.trade_log[0].close_reason, CloseReason::Target);
4834 assert_eq!(legacy.trade_log[0].size, future.trade_log[0].size);
4835 }
4836
4837 #[test]
4838 fn future_batch_sizes_from_shared_conversion_before_primary_and_uses_primary_eod() {
4839 let currency_plan = RunCurrencyPlan::new(
4840 "USD",
4841 ["EURUSD".to_owned()].into_iter().collect(),
4842 ["EURUSD".to_owned()].into_iter().collect(),
4843 [("EURUSD".to_owned(), "EUR".to_owned())]
4844 .into_iter()
4845 .collect(),
4846 [(
4847 "EUR".to_owned(),
4848 ConversionRoute::Direct {
4849 pair: FxPair {
4850 symbol: "EURUSD".to_owned(),
4851 base_currency: "EUR".to_owned(),
4852 quote_currency: "USD".to_owned(),
4853 },
4854 },
4855 )]
4856 .into_iter()
4857 .collect(),
4858 Vec::new(),
4859 )
4860 .unwrap();
4861 let mut config = fixed_lot_config();
4862 config.sizing = Some(SizingPolicy::FixedRiskAmount { amount: 12.0 });
4863 let future = FutureQuoteConfig {
4864 currency_plan: Some(currency_plan),
4865 conversion_stale_after_ms: 1_000,
4866 ..FutureQuoteConfig::default()
4867 };
4868 let events = vec![
4869 FeedEvent::new(
4870 tick("EURUSD", 1.1, 1.2, ts(10, 0, 0)),
4871 EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
4872 ),
4873 FeedEvent::new(
4874 tick("EURUSD", 2.0, 2.1, ts(10, 0, 1)),
4875 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
4876 ),
4877 ];
4878 let signals = vec![RawSignal::Entry {
4879 ts: ts(10, 0, 0),
4880 symbol: "EURUSD".into(),
4881 side: Side::Buy,
4882 order_type: OrderType::Market,
4883 price: Some(1.0),
4884 risk_multiplier: 1.0,
4885 stoploss: Some(1.19),
4886 targets: Vec::new(),
4887 group: None,
4888 trade_id: Some("shared-conversion".into()),
4889 }];
4890
4891 let mut feed = VecFeed::from_feed_events(events);
4892 let result = BacktestRunner::new_future(config, future)
4893 .run_raw_signals_future(&mut feed, signals, None);
4894
4895 assert_eq!(result.recorded_fills.len(), 2);
4896 assert!((result.recorded_fills[0].fill.price - 1.2).abs() < 1.0e-12);
4897 assert!((result.recorded_fills[0].size - 10.0).abs() < 1.0e-12);
4898 assert_eq!(result.recorded_fills[1].execution_ts, Some(ts(10, 0, 0)));
4899 assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 0));
4900 assert!(
4901 result
4902 .mtm_equity_curve
4903 .iter()
4904 .all(|point| point.ts == ts(10, 0, 0))
4905 );
4906 }
4907
4908 #[test]
4909 fn conversion_only_batch_revalues_but_defers_execution_to_primary_quote() {
4910 let currency_plan = RunCurrencyPlan::new(
4911 "USD",
4912 ["EURUSD".to_owned()].into_iter().collect(),
4913 ["EURUSD".to_owned()].into_iter().collect(),
4914 [("EURUSD".to_owned(), "EUR".to_owned())]
4915 .into_iter()
4916 .collect(),
4917 [(
4918 "EUR".to_owned(),
4919 ConversionRoute::Direct {
4920 pair: FxPair {
4921 symbol: "EURUSD".to_owned(),
4922 base_currency: "EUR".to_owned(),
4923 quote_currency: "USD".to_owned(),
4924 },
4925 },
4926 )]
4927 .into_iter()
4928 .collect(),
4929 Vec::new(),
4930 )
4931 .unwrap();
4932 let config = BacktestConfig {
4933 close_on_finish: false,
4934 ..fixed_lot_config()
4935 };
4936 let future = FutureQuoteConfig {
4937 currency_plan: Some(currency_plan),
4938 conversion_stale_after_ms: 10_000,
4939 ..FutureQuoteConfig::default()
4940 };
4941 let events = vec![
4942 FeedEvent::new(
4943 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
4944 EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
4945 ),
4946 FeedEvent::new(
4947 tick("EURUSD", 2.0, 2.0, ts(10, 0, 1)),
4948 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
4949 ),
4950 FeedEvent::new(
4951 tick("EURUSD", 110.0, 110.0, ts(10, 0, 2)),
4952 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
4953 ),
4954 ];
4955 let signals = vec![
4956 RawSignal::Entry {
4957 ts: ts(10, 0, 0),
4958 symbol: "EURUSD".into(),
4959 side: Side::Buy,
4960 order_type: OrderType::Market,
4961 price: None,
4962 risk_multiplier: 1.0,
4963 stoploss: None,
4964 targets: Vec::new(),
4965 group: None,
4966 trade_id: Some("conversion-only".into()),
4967 },
4968 RawSignal::Close {
4969 ts: ts(10, 0, 1),
4970 position: PositionRef::ByTradeId {
4971 trade_id: "conversion-only".into(),
4972 },
4973 },
4974 ];
4975
4976 let mut feed = VecFeed::from_feed_events(events);
4977 let result = BacktestRunner::new_future(config, future)
4978 .run_raw_signals_future(&mut feed, signals, None);
4979
4980 assert_eq!(result.recorded_fills.len(), 2);
4981 assert_eq!(result.recorded_fills[0].quote_ts, ts(10, 0, 0));
4982 assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 2));
4983 assert!(
4984 result
4985 .mtm_equity_curve
4986 .iter()
4987 .any(|point| point.ts == ts(10, 0, 1))
4988 );
4989 assert_eq!(result.total_pnl, 20.0);
4990 assert_eq!(result.close_events[0].native_pnl, Some(10.0));
4991 assert_eq!(
4992 result.close_events[0]
4993 .pnl_conversion
4994 .as_ref()
4995 .unwrap()
4996 .operation_ts,
4997 ts(10, 0, 2)
4998 );
4999 }
5000
5001 #[test]
5002 fn exact_timestamp_close_updates_balance_before_later_risk_entry() {
5003 let mut config = fixed_lot_config();
5004 config.close_on_finish = false;
5005 config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
5006 let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
5007 spec.digits = 2;
5008 spec.pip_position = 2;
5009 spec.lot_base_units = 1;
5010 spec.lot_step_units = 1;
5011 let future = FutureQuoteConfig {
5012 currency_plan: Some(identity_currency_plan("EURUSD")),
5013 ..FutureQuoteConfig::default()
5014 };
5015 let signals = vec![
5016 RawSignal::Entry {
5017 ts: ts(10, 0, 0),
5018 symbol: "EURUSD".into(),
5019 side: Side::Buy,
5020 order_type: OrderType::Market,
5021 price: None,
5022 risk_multiplier: 1.0,
5023 stoploss: Some(99.0),
5024 targets: Vec::new(),
5025 group: None,
5026 trade_id: Some("first".into()),
5027 },
5028 RawSignal::Close {
5029 ts: ts(10, 0, 1),
5030 position: PositionRef::ByTradeId {
5031 trade_id: "first".into(),
5032 },
5033 },
5034 RawSignal::Entry {
5035 ts: ts(10, 0, 1),
5036 symbol: "EURUSD".into(),
5037 side: Side::Buy,
5038 order_type: OrderType::Market,
5039 price: None,
5040 risk_multiplier: 1.0,
5041 stoploss: Some(100.0),
5042 targets: Vec::new(),
5043 group: None,
5044 trade_id: Some("second".into()),
5045 },
5046 ];
5047 let mut feed = VecFeed::new(vec![
5048 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
5049 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
5050 ]);
5051
5052 let result = BacktestRunner::new_future(config, future)
5053 .run_raw_signals_future(&mut feed, signals, None);
5054
5055 assert!((result.total_pnl - 100.0).abs() < 1.0e-12);
5056 assert_eq!(result.open_position_snapshots.len(), 1);
5057 assert_eq!(
5058 result.open_position_snapshots[0].trade_id.as_deref(),
5059 Some("second")
5060 );
5061 assert!((result.open_position_snapshots[0].remaining_size - 101.0).abs() < 1.0e-12);
5062 }
5063
5064 #[test]
5065 fn pending_fill_keeps_placement_size_after_balance_changes() {
5066 let mut config = fixed_lot_config();
5067 config.close_on_finish = false;
5068 config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
5069 let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
5070 spec.digits = 2;
5071 spec.pip_position = 2;
5072 spec.lot_base_units = 1;
5073 spec.lot_step_units = 1;
5074 let future = FutureQuoteConfig {
5075 currency_plan: Some(identity_currency_plan("EURUSD")),
5076 ..FutureQuoteConfig::default()
5077 };
5078 let signals = vec![
5079 RawSignal::Entry {
5080 ts: ts(10, 0, 0),
5081 symbol: "EURUSD".into(),
5082 side: Side::Buy,
5083 order_type: OrderType::Market,
5084 price: None,
5085 risk_multiplier: 1.0,
5086 stoploss: Some(99.0),
5087 targets: Vec::new(),
5088 group: None,
5089 trade_id: Some("market".into()),
5090 },
5091 RawSignal::Entry {
5092 ts: ts(10, 0, 0),
5093 symbol: "EURUSD".into(),
5094 side: Side::Buy,
5095 order_type: OrderType::Limit,
5096 price: Some(99.0),
5097 risk_multiplier: 1.0,
5098 stoploss: Some(98.0),
5099 targets: Vec::new(),
5100 group: None,
5101 trade_id: Some("pending".into()),
5102 },
5103 RawSignal::Close {
5104 ts: ts(10, 0, 1),
5105 position: PositionRef::ByTradeId {
5106 trade_id: "market".into(),
5107 },
5108 },
5109 ];
5110 let mut feed = VecFeed::new(vec![
5111 tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
5112 tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
5113 tick("EURUSD", 99.0, 99.0, ts(10, 0, 2)),
5114 ]);
5115
5116 let result = BacktestRunner::new_future(config, future)
5117 .run_raw_signals_future(&mut feed, signals, None);
5118
5119 assert_eq!(result.pending_order_snapshots.len(), 0);
5120 assert_eq!(result.open_position_snapshots.len(), 1);
5121 assert_eq!(
5122 result.open_position_snapshots[0].trade_id.as_deref(),
5123 Some("pending")
5124 );
5125 assert!((result.open_position_snapshots[0].remaining_size - 100.0).abs() < 1.0e-12);
5126 }
5127
5128 #[test]
5129 fn raw_entries_require_sizing_but_management_only_replay_does_not() {
5130 let entry = RawSignal::Entry {
5131 ts: ts(10, 0, 0),
5132 symbol: "EURUSD".into(),
5133 side: Side::Buy,
5134 order_type: OrderType::Market,
5135 price: None,
5136 risk_multiplier: 1.0,
5137 stoploss: None,
5138 targets: Vec::new(),
5139 group: None,
5140 trade_id: None,
5141 };
5142 let mut entry_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5143 let rejected =
5144 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
5145 .run_raw_signals_future(&mut entry_feed, vec![entry], None);
5146 assert!(rejected.action_dispositions.iter().any(|disposition| {
5147 disposition.action_id == "configuration"
5148 && disposition
5149 .reason
5150 .as_deref()
5151 .is_some_and(|reason| reason.contains("BacktestConfig.sizing"))
5152 }));
5153
5154 let mut management_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5155 let management =
5156 BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
5157 .run_raw_signals_future(
5158 &mut management_feed,
5159 vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
5160 None,
5161 );
5162 assert!(
5163 management
5164 .action_dispositions
5165 .iter()
5166 .all(|disposition| disposition.action_id != "configuration")
5167 );
5168 }
5169
5170 #[test]
5171 fn server_filter_signals_before_market_window() {
5172 let events = vec![
5178 tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
5179 tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
5180 ];
5181 let mut feed = VecFeed::new(events);
5182
5183 let raw_signals = vec![RawSignal::Entry {
5185 ts: NaiveDate::from_ymd_opt(2026, 1, 1)
5186 .unwrap()
5187 .and_hms_opt(0, 0, 0)
5188 .unwrap(),
5189 symbol: "EURUSD".into(),
5190 side: Side::Buy,
5191 order_type: OrderType::Market,
5192 price: Some(1.0850),
5193 risk_multiplier: 1.0,
5194 stoploss: None,
5195 targets: vec![1.0900],
5196 group: None,
5197 trade_id: None,
5198 }];
5199
5200 let runner = BacktestRunner::new(fixed_lot_config());
5201 let result = runner.run_raw_signals(&mut feed, raw_signals, None);
5202
5203 assert_eq!(result.total_trades, 1);
5205 }
5206}