Skip to main content

qs_backtest/
runner.rs

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