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