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