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