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