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