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::ScaleIn {
4297            position_id, size, ..
4298        } = &action
4299        {
4300            let valid = self
4301                .engine
4302                .get_position(position_id)
4303                .and_then(|position| explicit_instrument_spec(&self.config, &position.data.symbol))
4304                .is_none_or(|spec| {
4305                    size.to_string().parse::<Decimal>().is_ok_and(|quantity| {
4306                        quantity >= spec.quantity.minimum.get()
4307                            && spec
4308                                .quantity
4309                                .maximum
4310                                .is_none_or(|maximum| quantity <= maximum.get())
4311                            && spec.quantity.grid.contains(quantity).unwrap_or(false)
4312                    })
4313                });
4314            if !valid {
4315                let mut disposition =
4316                    ActionDisposition::rejected(action_id, "invalid_instrument_quantity");
4317                disposition.action_kind = Some(action_kind);
4318                disposition.signal_ts = Some(signal_ts);
4319                disposition.effective_ts = Some(effective_ts);
4320                self.record_disposition(lifecycle, disposition);
4321                return false;
4322            }
4323        }
4324        if let Action::Open {
4325            trade_id: Some(trade_id),
4326            ..
4327        } = &action
4328            && self.engine.manager.id_by_trade_id(trade_id).is_some()
4329        {
4330            let mut disposition = ActionDisposition::rejected(action_id, "duplicate_trade_id");
4331            disposition.action_kind = Some(action_kind);
4332            disposition.signal_ts = Some(signal_ts);
4333            disposition.effective_ts = Some(effective_ts);
4334            self.record_disposition(lifecycle, disposition);
4335            return false;
4336        }
4337        if let Action::ScaleIn { position_id, .. } = &action
4338            && future_executor.has_close(position_id)
4339        {
4340            let mut disposition =
4341                ActionDisposition::rejected(action_id, "scale_in_after_close_not_supported");
4342            disposition.action_kind = Some(action_kind);
4343            disposition.signal_ts = Some(signal_ts);
4344            disposition.effective_ts = Some(effective_ts);
4345            disposition.position_ids.push(position_id.clone());
4346            self.record_disposition(lifecycle, disposition);
4347            return false;
4348        }
4349
4350        let execution = match self.prepare_future_action(&mut action, execution, quote, pricer) {
4351            Ok(execution) => execution,
4352            Err(reason) => {
4353                let mut disposition = ActionDisposition::rejected(action_id, reason);
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 engine_transaction = match execution {
4363            Some(execution) => self
4364                .engine
4365                .begin_priced_future_action(action, quote, execution),
4366            None => self.engine.begin_future_action(action, effective_ts),
4367        };
4368        let engine_transaction = match engine_transaction {
4369            Ok(transaction) => transaction,
4370            Err(error) => {
4371                // A closed-position state mismatch on a management action mirrors
4372                // a live broker no-op, so it is skipped rather than failed.
4373                let closed_state = matches!(
4374                    error,
4375                    FutureApplyError::Core(qs_core::CoreError::InvalidState { .. })
4376                ) && action_kind != "entry";
4377                let mut disposition = if closed_state {
4378                    ActionDisposition::skipped(action_id, "position_closed")
4379                } else {
4380                    ActionDisposition::rejected(action_id, error.to_string())
4381                };
4382                disposition.action_kind = Some(action_kind);
4383                disposition.signal_ts = Some(signal_ts);
4384                disposition.effective_ts = Some(effective_ts);
4385                self.record_disposition(lifecycle, disposition);
4386                return false;
4387            }
4388        };
4389
4390        let committed_effects = engine_transaction.effects().to_vec();
4391        let mut affected = if FutureExecutor::requires_processing(engine_transaction.effects()) {
4392            match future_executor.process_future_effects_with_currency(
4393                engine_transaction.effects(),
4394                &self.engine,
4395                quote,
4396                Some(&action_id),
4397                Some(signal_ts),
4398                effective_ts,
4399                portfolio,
4400                Some(conversion_quotes),
4401            ) {
4402                Ok(affected) => affected,
4403                Err(error) => {
4404                    engine_transaction.rollback(&mut self.engine);
4405                    let mut disposition = ActionDisposition::failed(action_id, error.to_string());
4406                    disposition.action_kind = Some(action_kind);
4407                    disposition.signal_ts = Some(signal_ts);
4408                    disposition.effective_ts = Some(effective_ts);
4409                    self.record_disposition(lifecycle, disposition);
4410                    return false;
4411                }
4412            }
4413        } else {
4414            Vec::new()
4415        };
4416        for future_effect in engine_transaction.effects() {
4417            match future_effect.effect() {
4418                Effect::OrderPlaced { id } | Effect::OrderCancelled { id } => {
4419                    affected.push(id.clone());
4420                }
4421                _ => {}
4422            }
4423        }
4424        affected.sort();
4425        affected.dedup();
4426        let _ = engine_transaction.commit();
4427        self.record_committed_effects(committed_effects, Some(action_id.clone()));
4428
4429        let mut disposition = ActionDisposition::applied(action_id);
4430        disposition.action_kind = Some(action_kind);
4431        disposition.signal_ts = Some(signal_ts);
4432        disposition.effective_ts = Some(effective_ts);
4433        disposition.position_ids = affected;
4434        self.record_disposition(lifecycle, disposition);
4435        true
4436    }
4437
4438    fn prepare_future_action(
4439        &self,
4440        action: &mut Action,
4441        prepriced: Option<ExecutionFill>,
4442        quote: &PriceQuote,
4443        pricer: &ExecutionPricer,
4444    ) -> Result<Option<ExecutionFill>, String> {
4445        let mut execution = prepriced;
4446        match action {
4447            Action::Open {
4448                symbol,
4449                side,
4450                order_type,
4451                price,
4452                size,
4453                ..
4454            } => {
4455                if !valid_accounting_size(*size) {
4456                    return Err(format!(
4457                        "position size must be finite and greater than the accounting tolerance, got {size}"
4458                    ));
4459                }
4460                if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
4461                    return Err(format!(
4462                        "supplied entry price must be finite and positive, got {price:?}"
4463                    ));
4464                }
4465                if *order_type == OrderType::Market {
4466                    let priced = match execution {
4467                        Some(priced) => priced,
4468                        None => pricer
4469                            .market_entry(*side, quote, self.pip_size(symbol))
4470                            .map_err(|error| error.to_string())?,
4471                    };
4472                    *price = Some(priced.price);
4473                    execution = Some(priced);
4474                } else {
4475                    if price.is_none() {
4476                        return Err("pending entry requires a requested price".to_owned());
4477                    }
4478                    execution = None;
4479                }
4480            }
4481            Action::ScaleIn {
4482                position_id,
4483                price,
4484                size,
4485                ..
4486            } => {
4487                if !valid_accounting_size(*size) {
4488                    return Err(format!(
4489                        "scale-in size must be finite and greater than the accounting tolerance, got {size}"
4490                    ));
4491                }
4492                if price.is_some_and(|price| !price.is_finite() || price <= 0.0) {
4493                    return Err(format!(
4494                        "supplied scale-in price must be finite and positive, got {price:?}"
4495                    ));
4496                }
4497                let side = self
4498                    .engine
4499                    .get_position(position_id)
4500                    .map(|position| position.data.side)
4501                    .ok_or_else(|| format!("position not found: {position_id}"))?;
4502                let priced = match execution {
4503                    Some(priced) => priced,
4504                    None => pricer
4505                        .market_entry(side, quote, self.pip_size(&quote.symbol))
4506                        .map_err(|error| error.to_string())?,
4507                };
4508                *price = Some(priced.price);
4509                execution = Some(priced);
4510            }
4511            Action::ClosePosition { position_id } | Action::ClosePartial { position_id, .. } => {
4512                let position = self
4513                    .engine
4514                    .get_position(position_id)
4515                    .ok_or_else(|| format!("position not found: {position_id}"))?;
4516                if position.data.symbol != quote.symbol {
4517                    return Err(format!(
4518                        "position symbol {} does not match quote symbol {}",
4519                        position.data.symbol, quote.symbol
4520                    ));
4521                }
4522                execution = Some(
4523                    pricer
4524                        .market_exit(
4525                            position.data.side,
4526                            quote,
4527                            self.pip_size(&position.data.symbol),
4528                        )
4529                        .map_err(|error| error.to_string())?,
4530                );
4531            }
4532            _ => execution = None,
4533        }
4534        Ok(execution)
4535    }
4536
4537    fn prepare_triggering_pending(
4538        &self,
4539        quote: &PriceQuote,
4540        side: Option<Side>,
4541        pricer: &ExecutionPricer,
4542    ) -> (Vec<PreparedPendingFill>, Vec<(String, String)>) {
4543        let ids = self
4544            .engine
4545            .manager
4546            .pending_ids_by_symbol_sorted(&quote.symbol);
4547        let mut prepared = Vec::new();
4548        let mut failures = Vec::new();
4549        for id in ids {
4550            let Some(position) = self.engine.get_position(&id) else {
4551                continue;
4552            };
4553            if side.is_some_and(|side| position.data.side != side) {
4554                continue;
4555            }
4556            let Some(purpose) = position.pending_fill_purpose(quote, self.config.fill_model) else {
4557                continue;
4558            };
4559            let execution = match pricer.price(
4560                purpose,
4561                position.data.side,
4562                quote,
4563                position.data.pending_price,
4564                self.pip_size(&quote.symbol),
4565            ) {
4566                Ok(fill) => fill,
4567                Err(error) => {
4568                    failures.push((id, error.to_string()));
4569                    continue;
4570                }
4571            };
4572
4573            let size = position.data.size;
4574            if !valid_accounting_size(size) {
4575                failures.push((
4576                    id,
4577                    format!(
4578                        "pending size must be finite and greater than the accounting tolerance, got {size}"
4579                    ),
4580                ));
4581                continue;
4582            }
4583
4584            prepared.push(PreparedPendingFill {
4585                position_id: id,
4586                execution,
4587                size,
4588            });
4589        }
4590        (prepared, failures)
4591    }
4592
4593    fn pip_size(&self, symbol: &str) -> f64 {
4594        self.config
4595            .symbol_specs
4596            .get(symbol)
4597            .map(|spec| 10_f64.powi(-(spec.pip_position as i32)))
4598            .unwrap_or(0.0001)
4599    }
4600
4601    fn action_symbol(&self, action: &Action) -> Option<String> {
4602        match action {
4603            Action::Open { symbol, .. } => Some(symbol.clone()),
4604            Action::ClosePosition { position_id }
4605            | Action::ClosePartial { position_id, .. }
4606            | Action::ModifyStoploss { position_id, .. }
4607            | Action::MoveStoplossToEntry { position_id }
4608            | Action::AddTarget { position_id, .. }
4609            | Action::RemoveTarget { position_id, .. }
4610            | Action::ModifyTarget { position_id, .. }
4611            | Action::AddRule { position_id, .. }
4612            | Action::RemoveRule { position_id, .. }
4613            | Action::ScaleIn { position_id, .. }
4614            | Action::CancelPending { position_id } => self
4615                .engine
4616                .get_position(position_id)
4617                .map(|position| position.data.symbol.clone()),
4618            Action::CloseAllOf { symbol } | Action::ModifyAllStoploss { symbol, .. } => {
4619                Some(symbol.clone())
4620            }
4621            _ => None,
4622        }
4623    }
4624
4625    fn resolve_future_actions(&self, signal: &RawSignal) -> Vec<Action> {
4626        match signal {
4627            RawSignal::CloseAllOf { symbol, .. } => self
4628                .engine
4629                .manager
4630                .open_ids_by_symbol_sorted(symbol)
4631                .into_iter()
4632                .map(|position_id| Action::ClosePosition { position_id })
4633                .collect(),
4634            RawSignal::CloseAll { .. } => self
4635                .engine
4636                .manager
4637                .ids_by_status_sorted(PositionStatus::Open)
4638                .into_iter()
4639                .map(|position_id| Action::ClosePosition { position_id })
4640                .collect(),
4641            RawSignal::CloseAllInGroup { group_id, .. } => {
4642                let mut ids = self.engine.manager.open_ids_by_group(group_id);
4643                ids.sort();
4644                ids.into_iter()
4645                    .map(|position_id| Action::ClosePosition { position_id })
4646                    .collect()
4647            }
4648            RawSignal::CancelAllPending { .. } => self
4649                .engine
4650                .manager
4651                .ids_by_status_sorted(PositionStatus::Pending)
4652                .into_iter()
4653                .map(|position_id| Action::CancelPending { position_id })
4654                .collect(),
4655            _ => resolve_signal(signal, &self.engine),
4656        }
4657    }
4658}
4659
4660fn map_strategy_driver_error<FeedError, StrategyError>(
4661    error: StrategyDriverError<StrategyError>,
4662) -> StrategyReplayError<FeedError, StrategyError> {
4663    match error {
4664        StrategyDriverError::Series(error) => StrategyReplayError::Series(error),
4665        StrategyDriverError::SeriesView(error) => StrategyReplayError::SeriesView(error),
4666        StrategyDriverError::Analysis(error) => StrategyReplayError::Analysis(error),
4667        StrategyDriverError::Strategy(error) => StrategyReplayError::Strategy(error),
4668        StrategyDriverError::Runtime(error) => StrategyReplayError::Runtime(error),
4669        StrategyDriverError::WarmupSignals { timestamp } => {
4670            StrategyReplayError::WarmupSignals { timestamp }
4671        }
4672        StrategyDriverError::InvalidGeneratedSignal {
4673            signal_index,
4674            reason,
4675        } => StrategyReplayError::InvalidGeneratedSignal {
4676            signal_index,
4677            reason,
4678        },
4679        StrategyDriverError::TickExecutionRequired { symbol, timestamp } => {
4680            StrategyReplayError::TickExecutionRequired { symbol, timestamp }
4681        }
4682    }
4683}
4684
4685#[allow(clippy::too_many_arguments)]
4686fn observe_future_equity(
4687    portfolio: &mut PortfolioRecorder,
4688    future_executor: &FutureExecutor,
4689    ts: NaiveDateTime,
4690    conversion_quotes: &ConversionQuoteBook,
4691    kind: EquityObservationKind,
4692    collector: &mut MtmCurveCollector,
4693    last_candidate: &mut Option<EquityPoint>,
4694    suppress_unchanged_post_output: bool,
4695) {
4696    portfolio.set_realized_pnl(future_executor.realized_pnl());
4697    let mut point = portfolio.observe_with_currency(
4698        ts,
4699        future_executor.open_snapshots(),
4700        Some(conversion_quotes),
4701    );
4702    point.observation_kind = Some(kind.as_str().to_owned());
4703    if suppress_unchanged_post_output
4704        && last_candidate
4705            .as_ref()
4706            .is_some_and(|previous| same_equity_values(previous, &point))
4707    {
4708        return;
4709    }
4710    collector.observe(point.clone());
4711    *last_candidate = Some(point);
4712}
4713
4714fn same_equity_values(left: &EquityPoint, right: &EquityPoint) -> bool {
4715    let mut left = left.clone();
4716    let mut right = right.clone();
4717    left.observation_kind = None;
4718    left.observation_sequence = None;
4719    right.observation_kind = None;
4720    right.observation_sequence = None;
4721    left == right
4722}
4723
4724fn pending_fill_position_id(effect: &FutureEffect) -> Option<&str> {
4725    match effect.effect() {
4726        Effect::PositionOpened { id } => Some(id),
4727        _ => None,
4728    }
4729}
4730
4731fn valid_accounting_size(size: f64) -> bool {
4732    size.is_finite() && size > position_size_tolerance(size)
4733}
4734
4735fn explicit_instrument_spec<'a>(
4736    config: &'a BacktestConfig,
4737    symbol: &str,
4738) -> Option<&'a InstrumentSpec> {
4739    config
4740        .instrument_manifest
4741        .as_ref()?
4742        .instruments
4743        .get(symbol)
4744        .map(|artifact| &artifact.spec)
4745}
4746
4747fn decimal_to_f64(value: Decimal, field: &str) -> Result<f64, String> {
4748    let value = value
4749        .to_string()
4750        .parse::<f64>()
4751        .map_err(|error| format!("invalid {field}: {error}"))?;
4752    if value.is_finite() {
4753        Ok(value)
4754    } else {
4755        Err(format!("{field} must be finite"))
4756    }
4757}
4758
4759fn instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
4760    decimal_to_f64(
4761        spec.economics.contract_multiplier.get(),
4762        "instrument contract multiplier",
4763    )
4764    .and_then(|value| {
4765        if value > 0.0 {
4766            Ok(value)
4767        } else {
4768            Err("instrument contract multiplier must be positive".into())
4769        }
4770    })
4771}
4772
4773fn supported_instrument_multiplier(spec: &InstrumentSpec) -> Result<f64, String> {
4774    if spec.status != ListingStatus::Trading {
4775        return Err(format!(
4776            "instrument {} is not in trading status",
4777            spec.instrument
4778        ));
4779    }
4780    if spec.economics.quantity_unit != QuantityUnit::StandardLot {
4781        return Err(format!(
4782            "unsupported quantity unit for instrument {}: {:?}",
4783            spec.instrument, spec.economics.quantity_unit
4784        ));
4785    }
4786    let model = spec.economics.pnl_model.as_str();
4787    if model != EconomicsModelId::FX_QUOTE_LINEAR_V1
4788        && model != EconomicsModelId::CFD_QUOTE_LINEAR_V1
4789    {
4790        return Err(format!(
4791            "unsupported P&L model for instrument {}: {model}",
4792            spec.instrument
4793        ));
4794    }
4795    instrument_multiplier(spec)
4796}
4797
4798fn validate_instrument_manifest(config: &BacktestConfig) -> Result<(), String> {
4799    let Some(manifest) = &config.instrument_manifest else {
4800        return Ok(());
4801    };
4802    for (symbol, artifact) in &manifest.instruments {
4803        if symbol.is_empty() {
4804            return Err("instrument manifest symbol must not be empty".into());
4805        }
4806        artifact
4807            .spec
4808            .validate()
4809            .map_err(|error| format!("invalid instrument spec for {symbol}: {error}"))?;
4810        if artifact.resolved.instrument != artifact.spec.instrument {
4811            return Err(format!(
4812                "resolved instrument and spec identity differ for {symbol}"
4813            ));
4814        }
4815        if artifact.resolved.spec_revision != artifact.spec.revision {
4816            return Err(format!(
4817                "resolved specification revision does not match the embedded spec for {symbol}"
4818            ));
4819        }
4820        supported_instrument_multiplier(&artifact.spec)?;
4821    }
4822    for binding in &manifest.stored_series {
4823        let known = manifest
4824            .instruments
4825            .values()
4826            .any(|artifact| artifact.resolved == binding.instrument);
4827        if !known {
4828            return Err(format!(
4829                "stored series {}:{} references an instrument outside the manifest",
4830                binding.source_partition, binding.source_symbol
4831            ));
4832        }
4833        let artifact = manifest
4834            .instruments
4835            .values()
4836            .find(|artifact| artifact.resolved == binding.instrument)
4837            .expect("known binding reference has an instrument artifact");
4838        if binding.effective != artifact.spec.effective {
4839            return Err(format!(
4840                "stored series {}:{} effective interval differs from its instrument spec",
4841                binding.source_partition, binding.source_symbol
4842            ));
4843        }
4844    }
4845    Ok(())
4846}
4847
4848fn effective_contract_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
4849    let mut contract_sizes = config.contract_sizes.clone();
4850    if let Some(manifest) = &config.instrument_manifest {
4851        for (symbol, artifact) in &manifest.instruments {
4852            if let Ok(multiplier) = instrument_multiplier(&artifact.spec) {
4853                contract_sizes.insert(symbol.clone(), multiplier);
4854            }
4855        }
4856    }
4857    contract_sizes
4858}
4859
4860/// Reject cost specifications that cannot be applied deterministically to this run.
4861fn validate_replay_costs(
4862    config: &BacktestConfig,
4863    future: Option<&FutureQuoteConfig>,
4864) -> Result<(), String> {
4865    let account_currency = future
4866        .and_then(|future| future.currency_plan.as_ref())
4867        .map(|plan| plan.account_currency().to_owned());
4868    for (symbol, costs) in &config.costs {
4869        if symbol.is_empty() {
4870            return Err("cost symbol must not be empty".into());
4871        }
4872        costs
4873            .validate()
4874            .map_err(|error| format!("costs for {symbol} are invalid: {error}"))?;
4875        if let Some(account_currency) = account_currency.as_deref() {
4876            costs
4877                .validate_against_account_currency(account_currency)
4878                .map_err(|error| format!("costs for {symbol} are invalid: {error}"))?;
4879        }
4880        if costs.requires_point_size() && !config.symbol_specs.contains_key(symbol) {
4881            return Err(format!(
4882                "point-denominated swap for {symbol} requires a symbol specification for its digit count"
4883            ));
4884        }
4885    }
4886    Ok(())
4887}
4888
4889/// Price point size per symbol, derived from the digit count used by the symbol registry.
4890fn effective_point_sizes(config: &BacktestConfig) -> HashMap<String, f64> {
4891    config
4892        .symbol_specs
4893        .iter()
4894        .map(|(symbol, spec)| (symbol.clone(), 10f64.powi(-i32::from(spec.digits))))
4895        .collect()
4896}
4897
4898fn is_monetary_sizing(policy: &SizingPolicy) -> bool {
4899    matches!(
4900        policy,
4901        SizingPolicy::FixedRiskAmount { .. } | SizingPolicy::BalanceRiskPercent { .. }
4902    )
4903}
4904
4905fn accept_legacy_quote(
4906    quote: &PriceQuote,
4907    last_quote_ts: &mut BTreeMap<String, NaiveDateTime>,
4908) -> bool {
4909    if ExecutionPricer::validate_quote(quote).is_err()
4910        || last_quote_ts
4911            .get(&quote.symbol)
4912            .is_some_and(|last| *last > quote.ts)
4913    {
4914        return false;
4915    }
4916    last_quote_ts.insert(quote.symbol.clone(), quote.ts);
4917    true
4918}
4919
4920/// Largest number of run tags one replay may carry.
4921pub const MAX_RUN_TAGS: usize = 32;
4922/// Largest byte length of one run-tag key or value.
4923pub const MAX_RUN_TAG_BYTES: usize = 64;
4924
4925fn validate_run_tags(tags: &BTreeMap<String, String>) -> Result<(), String> {
4926    if tags.len() > MAX_RUN_TAGS {
4927        return Err(format!(
4928            "run tags must not exceed {MAX_RUN_TAGS} entries, got {}",
4929            tags.len()
4930        ));
4931    }
4932    for (key, value) in tags {
4933        if key.is_empty() || key.len() > MAX_RUN_TAG_BYTES {
4934            return Err(format!(
4935                "run tag key must be 1 to {MAX_RUN_TAG_BYTES} bytes, got '{key}'"
4936            ));
4937        }
4938        if !key
4939            .chars()
4940            .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
4941        {
4942            return Err(format!(
4943                "run tag key must be ASCII alphanumeric, underscore, or hyphen, got '{key}'"
4944            ));
4945        }
4946        if value.len() > MAX_RUN_TAG_BYTES {
4947            return Err(format!(
4948                "run tag value for '{key}' must not exceed {MAX_RUN_TAG_BYTES} bytes"
4949            ));
4950        }
4951        if value.chars().any(char::is_control) {
4952            return Err(format!(
4953                "run tag value for '{key}' must not contain control characters"
4954            ));
4955        }
4956    }
4957    Ok(())
4958}
4959
4960fn validate_replay_config(
4961    config: &BacktestConfig,
4962    future: Option<&FutureQuoteConfig>,
4963    raw_signals: &[RawSignal],
4964) -> Result<(), String> {
4965    if !config.initial_balance.is_finite() || config.initial_balance <= 0.0 {
4966        return Err(format!(
4967            "initial balance must be finite and positive, got {}",
4968            config.initial_balance
4969        ));
4970    }
4971    for (symbol, contract_size) in &config.contract_sizes {
4972        if symbol.is_empty() {
4973            return Err("contract-size symbol must not be empty".into());
4974        }
4975        if !contract_size.is_finite() || *contract_size <= 0.0 {
4976            return Err(format!(
4977                "contract size for {symbol} must be finite and positive, got {contract_size}"
4978            ));
4979        }
4980    }
4981
4982    validate_run_tags(&config.run_tags)?;
4983    validate_replay_costs(config, future)?;
4984    validate_instrument_manifest(config)?;
4985    for (symbol, spec) in &config.symbol_specs {
4986        validate_symbol_spec(symbol, spec)?;
4987        if explicit_instrument_spec(config, symbol).is_none() {
4988            resolve_legacy_economics(spec).map_err(|error| error.to_string())?;
4989        }
4990    }
4991
4992    let entry_symbols: Vec<&str> = raw_signals
4993        .iter()
4994        .filter_map(|signal| match signal {
4995            RawSignal::Entry { symbol, .. } => Some(symbol.as_str()),
4996            _ => None,
4997        })
4998        .collect();
4999    if !entry_symbols.is_empty() && config.sizing.is_none() {
5000        return Err("raw entry requires BacktestConfig.sizing".to_owned());
5001    }
5002    if let Some(policy) = &config.sizing {
5003        validate_sizing_policy(policy)?;
5004        for symbol in &entry_symbols {
5005            if !config.symbol_specs.contains_key(*symbol)
5006                && explicit_instrument_spec(config, symbol).is_none()
5007            {
5008                return Err(format!("missing instrument or symbol spec for {symbol}"));
5009            }
5010        }
5011        if !entry_symbols.is_empty() && is_monetary_sizing(policy) {
5012            let future = future.ok_or_else(|| {
5013                "monetary sizing requires FutureQuote execution and a currency plan".to_owned()
5014            })?;
5015            let plan = future
5016                .currency_plan
5017                .as_ref()
5018                .ok_or_else(|| "monetary sizing requires a FutureQuote currency plan".to_owned())?;
5019            for symbol in &entry_symbols {
5020                if plan.route_for_primary_symbol(symbol).is_none() {
5021                    return Err(format!(
5022                        "currency plan has no frozen route for primary symbol {symbol}"
5023                    ));
5024                }
5025            }
5026        }
5027    }
5028
5029    if let Some(future) = future {
5030        if future.signal_latency_ms < 0 {
5031            return Err(format!(
5032                "signal latency must be non-negative, got {}",
5033                future.signal_latency_ms
5034            ));
5035        }
5036        let latency = Duration::milliseconds(future.signal_latency_ms);
5037        for signal in raw_signals {
5038            if signal.ts().checked_add_signed(latency).is_none() {
5039                return Err(format!(
5040                    "signal latency overflows datetime for signal at {}",
5041                    signal.ts()
5042                ));
5043            }
5044        }
5045        if !future.slippage_pips.is_finite() {
5046            return Err(format!(
5047                "slippage pips must be finite, got {}",
5048                future.slippage_pips
5049            ));
5050        }
5051        if future.stale_quote_after_ms.is_some_and(|value| value < 0) {
5052            return Err("stale quote threshold must be non-negative".into());
5053        }
5054        if !future.pnl_epsilon.is_finite() || future.pnl_epsilon < 0.0 {
5055            return Err(format!(
5056                "P&L epsilon must be finite and non-negative, got {}",
5057                future.pnl_epsilon
5058            ));
5059        }
5060        if future.conversion_stale_after_ms < 0 {
5061            return Err("conversion quote threshold must be non-negative".to_owned());
5062        }
5063        future
5064            .mtm_output
5065            .validate()
5066            .map_err(|error| error.to_string())?;
5067    }
5068    Ok(())
5069}
5070
5071fn validate_sizing_policy(policy: &SizingPolicy) -> Result<(), String> {
5072    let (name, value) = match policy {
5073        SizingPolicy::FixedLot { lots } => ("fixed lots", *lots),
5074        SizingPolicy::FixedRiskAmount { amount } => ("fixed risk amount", *amount),
5075        SizingPolicy::BalanceRiskPercent { percent } => ("balance risk percent", *percent),
5076    };
5077    if value.is_finite() && value > 0.0 {
5078        Ok(())
5079    } else {
5080        Err(format!("{name} must be finite and positive, got {value}"))
5081    }
5082}
5083
5084fn validate_symbol_spec(symbol: &str, spec: &qs_symbols::SymbolSpec) -> Result<(), String> {
5085    if symbol.is_empty() || spec.canonical.is_empty() {
5086        return Err("symbol spec names must not be empty".into());
5087    }
5088    if spec.digits > 18 || spec.pip_position > spec.digits {
5089        return Err(format!(
5090            "invalid price precision for {symbol}: digits={}, pip_position={}",
5091            spec.digits, spec.pip_position
5092        ));
5093    }
5094    if spec.lot_base_units <= 0
5095        || spec.lot_step_units <= 0
5096        || spec.lot_min_steps <= 0
5097        || spec.lot_max_steps < 0
5098        || (spec.lot_max_steps > 0 && spec.lot_max_steps < spec.lot_min_steps)
5099    {
5100        return Err(format!("invalid lot metadata for {symbol}"));
5101    }
5102    let lot_step = spec.lot_step();
5103    let min_lot = spec.lot_min();
5104    let max_lot = spec.lot_max();
5105    if !lot_step.is_finite()
5106        || lot_step <= 0.0
5107        || !min_lot.is_finite()
5108        || min_lot <= 0.0
5109        || !max_lot.is_finite()
5110    {
5111        return Err(format!("invalid derived lot metadata for {symbol}"));
5112    }
5113    Ok(())
5114}
5115
5116fn rejected_legacy_result(config: &BacktestConfig) -> BacktestResult {
5117    BacktestResult::from_trade_log(
5118        if config.initial_balance.is_finite() {
5119            config.initial_balance
5120        } else {
5121            0.0
5122        },
5123        Vec::new(),
5124    )
5125}
5126
5127fn rejected_future_result(
5128    config: &BacktestConfig,
5129    future: &FutureQuoteConfig,
5130    evaluation_options: EvaluationOptions,
5131    error: String,
5132) -> BacktestResult {
5133    let execution_model = ExecutionModel::new(
5134        qs_core::types::ExecutionConvention::FutureQuoteV1,
5135        config.fill_model,
5136        if future.slippage_pips == 0.0 {
5137            SlippageModel::None
5138        } else {
5139            SlippageModel::FixedPips {
5140                pips: future.slippage_pips,
5141            }
5142        },
5143    );
5144    let mut lifecycle = LifecycleLedger::new();
5145    let _ = lifecycle.record(ActionDisposition::rejected(
5146        "configuration",
5147        format!("invalid_configuration: {error}"),
5148    ));
5149    let mut tags = BTreeMap::new();
5150    tags.insert("configuration_error".into(), error);
5151    insert_economic_support_metadata(&mut tags, config);
5152    let artifacts = FutureBacktestArtifacts {
5153        execution: ExecutionMetadata {
5154            execution_model,
5155            initial_balance: if config.initial_balance.is_finite() {
5156                config.initial_balance
5157            } else {
5158                0.0
5159            },
5160            account_currency: future
5161                .currency_plan
5162                .as_ref()
5163                .map(|plan| plan.account_currency().to_owned()),
5164            currency_plan: future.currency_plan.clone(),
5165            contract_sizes: effective_contract_sizes(config)
5166                .into_iter()
5167                .filter(|(symbol, size)| !symbol.is_empty() && size.is_finite() && *size > 0.0)
5168                .collect(),
5169            instrument_manifest: config.instrument_manifest.clone(),
5170            instrument_sizing: Vec::new(),
5171            market_entry_sizing_basis: future.market_entry_sizing_basis,
5172            market_entry_sizing: Vec::new(),
5173            stale_quote_after_millis: future.stale_quote_after_ms,
5174            pnl_epsilon: if future.pnl_epsilon.is_finite() && future.pnl_epsilon >= 0.0 {
5175                future.pnl_epsilon
5176            } else {
5177                crate::artifacts::DEFAULT_PNL_EPSILON
5178            },
5179            tags,
5180            ..ExecutionMetadata::default()
5181        },
5182        lifecycle,
5183        mtm_output_summary: MtmOutputSummary {
5184            policy: future.mtm_output,
5185            ..MtmOutputSummary::default()
5186        },
5187        ..FutureBacktestArtifacts::default()
5188    };
5189    BacktestResult::from_future_artifacts_with_options(artifacts, evaluation_options)
5190}
5191
5192fn insert_economic_support_metadata(tags: &mut BTreeMap<String, String>, config: &BacktestConfig) {
5193    let mut compatibility_specs = config.symbol_specs.iter().peekable();
5194    if compatibility_specs.peek().is_none() {
5195        return;
5196    }
5197    tags.insert(
5198        "economics.guard".into(),
5199        LEGACY_ECONOMIC_GUARD_ID.to_owned(),
5200    );
5201    for (symbol, spec) in compatibility_specs {
5202        let prefix = format!("economics.symbol.{symbol}");
5203        tags.insert(format!("{prefix}.category"), spec.category.clone());
5204        match resolve_legacy_economics(spec) {
5205            Ok(economics) => {
5206                tags.insert(format!("{prefix}.status"), "supported".into());
5207                tags.insert(format!("{prefix}.model"), economics.model.as_str().into());
5208                tags.insert(
5209                    format!("{prefix}.contract_multiplier"),
5210                    economics.contract_multiplier.to_string(),
5211                );
5212            }
5213            Err(error) => {
5214                tags.insert(format!("{prefix}.status"), "unsupported".into());
5215                tags.insert(format!("{prefix}.reason"), error.to_string());
5216            }
5217        }
5218    }
5219}
5220
5221fn queued_exposure_symbols(
5222    queued: &VecDeque<QueuedAction>,
5223    quotes: &BTreeMap<String, PriceQuote>,
5224    batch_ts: NaiveDateTime,
5225) -> BTreeSet<String> {
5226    queued
5227        .iter()
5228        .filter(|action| {
5229            action.effective_ts <= batch_ts
5230                && quotes.contains_key(&action.symbol)
5231                && is_exposure_increasing(&action.action)
5232        })
5233        .map(|action| action.symbol.clone())
5234        .collect()
5235}
5236
5237fn is_exposure_increasing(action: &Action) -> bool {
5238    matches!(
5239        action,
5240        Action::Open {
5241            order_type: OrderType::Market,
5242            ..
5243        } | Action::ScaleIn { .. }
5244    )
5245}
5246
5247fn is_fill_bearing(action: &Action) -> bool {
5248    matches!(
5249        action,
5250        Action::Open {
5251            order_type: OrderType::Market,
5252            ..
5253        } | Action::ClosePosition { .. }
5254            | Action::ClosePartial { .. }
5255            | Action::ScaleIn { .. }
5256    )
5257}
5258
5259fn raw_signal_kind(signal: &RawSignal) -> &'static str {
5260    match signal {
5261        RawSignal::Entry { .. } => "entry",
5262        RawSignal::Close { .. } => "close",
5263        RawSignal::ClosePartial { .. } => "close_partial",
5264        RawSignal::ModifyStoploss { .. } => "modify_stoploss",
5265        RawSignal::MoveStoplossToEntry { .. } => "move_stoploss_to_entry",
5266        RawSignal::AddTarget { .. } => "add_target",
5267        RawSignal::RemoveTarget { .. } => "remove_target",
5268        RawSignal::ModifyTarget { .. } => "modify_target",
5269        RawSignal::AddRule { .. } => "add_rule",
5270        RawSignal::RemoveRule { .. } => "remove_rule",
5271        RawSignal::ScaleIn { .. } => "scale_in",
5272        RawSignal::CancelPending { .. } => "cancel_pending",
5273        RawSignal::CloseAllOf { .. } => "close_all_of",
5274        RawSignal::CloseAll { .. } => "close_all",
5275        RawSignal::CancelAllPending { .. } => "cancel_all_pending",
5276        RawSignal::ModifyAllStoploss { .. } => "modify_all_stoploss",
5277        RawSignal::CloseAllInGroup { .. } => "close_all_in_group",
5278        RawSignal::ModifyAllStoplossInGroup { .. } => "modify_all_stoploss_in_group",
5279    }
5280}
5281
5282// ─── Tests ──────────────────────────────────────────────────────────────────
5283
5284#[cfg(test)]
5285mod tests {
5286    use super::*;
5287    use crate::currency::{ConversionRoute, FxPair};
5288    use crate::data_feed::{EventMetadata, FeedEvent, MarketEvent, SeriesRoles, VecFeed};
5289    use crate::profile::{
5290        EntryGeometryPolicy, ManagementProfile, PositionRef, RawSignal, StoplossMode, TargetSource,
5291    };
5292    use chrono::NaiveDate;
5293    use qs_core::types::{CloseReason, FillPurpose, OrderType, Side, TargetSpec};
5294
5295    fn ts(h: u32, m: u32, s: u32) -> chrono::NaiveDateTime {
5296        NaiveDate::from_ymd_opt(2026, 1, 1)
5297            .unwrap()
5298            .and_hms_opt(h, m, s)
5299            .unwrap()
5300    }
5301
5302    fn tick(symbol: &str, bid: f64, ask: f64, time: chrono::NaiveDateTime) -> MarketEvent {
5303        MarketEvent::Tick {
5304            symbol: symbol.into(),
5305            ts: time,
5306            bid,
5307            ask,
5308        }
5309    }
5310
5311    #[test]
5312    fn scheduled_signal_preserves_an_explicit_opaque_action_id() {
5313        let scheduled = ScheduledSignal::new(
5314            7,
5315            ts(10, 0, 0),
5316            ts(10, 0, 1),
5317            RawSignal::CloseAll { ts: ts(10, 0, 0) },
5318            true,
5319        )
5320        .with_action_id("caller-command/opaque:7");
5321
5322        assert_eq!(scheduled.resolved_action_id(), "caller-command/opaque:7");
5323    }
5324
5325    #[test]
5326    fn scheduled_signal_keeps_the_compatible_generated_action_id() {
5327        let scheduled = ScheduledSignal::new(
5328            7,
5329            ts(10, 0, 0),
5330            ts(10, 0, 1),
5331            RawSignal::CloseAll { ts: ts(10, 0, 0) },
5332            false,
5333        );
5334
5335        assert_eq!(scheduled.resolved_action_id(), "signal:00000007");
5336    }
5337
5338    fn test_symbol_spec(symbol: &str) -> qs_symbols::SymbolSpec {
5339        qs_symbols::SymbolSpec {
5340            canonical: symbol.to_ascii_lowercase(),
5341            pip_position: 4,
5342            digits: 5,
5343            category: "forex".into(),
5344            lot_base_units: 100,
5345            lot_step_units: 1,
5346            lot_min_steps: 1,
5347            lot_max_steps: 0,
5348        }
5349    }
5350
5351    fn fixed_lot_config() -> BacktestConfig {
5352        BacktestConfig {
5353            sizing: Some(SizingPolicy::FixedLot { lots: 1.0 }),
5354            symbol_specs: ["EURUSD", "XAUUSD"]
5355                .into_iter()
5356                .map(|symbol| (symbol.to_owned(), test_symbol_spec(symbol)))
5357                .collect(),
5358            ..BacktestConfig::default()
5359        }
5360    }
5361
5362    fn identity_currency_plan(symbol: &str) -> RunCurrencyPlan {
5363        RunCurrencyPlan::new(
5364            "USD",
5365            [symbol.to_owned()].into_iter().collect(),
5366            Default::default(),
5367            [(symbol.to_owned(), "USD".to_owned())]
5368                .into_iter()
5369                .collect(),
5370            [(
5371                "USD".to_owned(),
5372                ConversionRoute::Identity {
5373                    currency: "USD".to_owned(),
5374                },
5375            )]
5376            .into_iter()
5377            .collect(),
5378            Vec::new(),
5379        )
5380        .unwrap()
5381    }
5382
5383    struct ScriptedBatchFeed {
5384        batches: VecDeque<Result<Option<TimestampBatch>, &'static str>>,
5385    }
5386
5387    impl FallibleBatchFeed for ScriptedBatchFeed {
5388        type Error = &'static str;
5389
5390        fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
5391            self.batches.pop_front().unwrap_or(Ok(None))
5392        }
5393    }
5394
5395    struct CountingBatchFeed {
5396        batches: VecDeque<TimestampBatch>,
5397        polls: std::rc::Rc<std::cell::Cell<usize>>,
5398    }
5399
5400    impl FallibleBatchFeed for CountingBatchFeed {
5401        type Error = Infallible;
5402
5403        fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
5404            self.polls.set(self.polls.get() + 1);
5405            Ok(self.batches.pop_front())
5406        }
5407    }
5408
5409    fn primary_batch(event: MarketEvent) -> TimestampBatch {
5410        TimestampBatch {
5411            ts: event.ts(),
5412            events: vec![FeedEvent::new(
5413                event,
5414                EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
5415            )],
5416        }
5417    }
5418
5419    fn market_entry(timestamp: NaiveDateTime, symbol: &str, order_type: OrderType) -> RawSignal {
5420        RawSignal::Entry {
5421            ts: timestamp,
5422            symbol: symbol.into(),
5423            side: Side::Buy,
5424            order_type,
5425            price: (order_type == OrderType::Limit).then_some(1.0),
5426            risk_multiplier: 1.0,
5427            stoploss: None,
5428            targets: Vec::new(),
5429            group: None,
5430            trade_id: Some(format!("{symbol}-blocker")),
5431            entry_class: None,
5432        }
5433    }
5434
5435    #[test]
5436    fn future_streaming_matches_materialized_and_stops_without_draining() {
5437        let events = vec![
5438            tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5439            tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
5440            tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
5441        ];
5442        let signals = vec![
5443            market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market),
5444            RawSignal::CloseAll { ts: ts(10, 0, 1) },
5445        ];
5446        let config = BacktestConfig {
5447            close_on_finish: false,
5448            ..fixed_lot_config()
5449        };
5450        let mut materialized_feed = VecFeed::new(events.clone());
5451        let materialized = BacktestRunner::new_future(
5452            config.clone(),
5453            FutureQuoteConfig {
5454                mtm_output: MtmOutputPolicy::Full,
5455                ..FutureQuoteConfig::default()
5456            },
5457        )
5458        .run_raw_signals_future(&mut materialized_feed, signals.clone(), None);
5459
5460        let mut stream = ScriptedBatchFeed {
5461            batches: VecDeque::from([
5462                Ok(Some(primary_batch(events[0].clone()))),
5463                Ok(Some(primary_batch(events[1].clone()))),
5464                Ok(Some(primary_batch(events[2].clone()))),
5465                Err("must not drain"),
5466            ]),
5467        };
5468        let mut progress = Vec::new();
5469        let streamed = BacktestRunner::new_future(
5470            config,
5471            FutureQuoteConfig {
5472                mtm_output: MtmOutputPolicy::Full,
5473                ..FutureQuoteConfig::default()
5474            },
5475        )
5476        .run_raw_signals_future_streaming_controlled(
5477            &mut stream,
5478            Some(ts(10, 0, 2)),
5479            signals,
5480            None,
5481            || false,
5482            |update| progress.push(update),
5483        )
5484        .unwrap();
5485
5486        assert_eq!(
5487            serde_json::to_value(&streamed).unwrap(),
5488            serde_json::to_value(&materialized).unwrap()
5489        );
5490        assert_eq!(
5491            stream.batches.len(),
5492            2,
5493            "quiescence must leave the tail unread"
5494        );
5495        assert_eq!(progress.first().unwrap().total_events, 0);
5496        assert_eq!(progress.last().unwrap().processed_events, 2);
5497        assert_eq!(progress.last().unwrap().total_events, 2);
5498        assert_eq!(
5499            streamed.mtm_equity_curve.last().unwrap().ts,
5500            ts(10, 0, 1),
5501            "terminal observation must use the last processed primary timestamp"
5502        );
5503        assert_eq!(
5504            streamed
5505                .mtm_equity_curve
5506                .last()
5507                .unwrap()
5508                .observation_kind
5509                .as_deref(),
5510            Some(EquityObservationKind::QuiescentTermination.as_str())
5511        );
5512        assert_eq!(
5513            streamed
5514                .execution_metadata
5515                .as_ref()
5516                .unwrap()
5517                .tags
5518                .get("termination_reason")
5519                .map(String::as_str),
5520            Some("quiescent")
5521        );
5522    }
5523
5524    #[test]
5525    fn exact_time_close_waits_for_later_symbol_pending_fill() {
5526        let open_ts = ts(10, 0, 0);
5527        let execution_ts = ts(10, 0, 1);
5528        let events = vec![
5529            FeedEvent::new(
5530                tick("XAUUSD", 101.0, 101.0, open_ts),
5531                EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
5532            ),
5533            FeedEvent::new(
5534                tick("EURUSD", 1.1, 1.1, execution_ts),
5535                EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
5536            ),
5537            FeedEvent::new(
5538                tick("XAUUSD", 100.0, 100.0, execution_ts),
5539                EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
5540            ),
5541        ];
5542        let signals = vec![
5543            RawSignal::Entry {
5544                ts: open_ts,
5545                symbol: "XAUUSD".into(),
5546                side: Side::Buy,
5547                order_type: OrderType::Limit,
5548                price: Some(100.0),
5549                risk_multiplier: 1.0,
5550                stoploss: None,
5551                targets: Vec::new(),
5552                group: None,
5553                trade_id: Some("later-pending".into()),
5554                entry_class: None,
5555            },
5556            RawSignal::Close {
5557                ts: execution_ts,
5558                position: PositionRef::ByTradeId {
5559                    trade_id: "later-pending".into(),
5560                },
5561            },
5562        ];
5563        let mut feed = VecFeed::from_feed_events(events);
5564        let result = BacktestRunner::new_future(
5565            BacktestConfig {
5566                close_on_finish: false,
5567                ..fixed_lot_config()
5568            },
5569            FutureQuoteConfig::default(),
5570        )
5571        .run_raw_signals_future(&mut feed, signals, None);
5572
5573        assert_eq!(
5574            result
5575                .recorded_fills
5576                .iter()
5577                .map(|fill| fill.fill.purpose)
5578                .collect::<Vec<_>>(),
5579            vec![FillPurpose::LimitEntry, FillPurpose::MarketExit]
5580        );
5581        assert_eq!(result.close_events.len(), 1);
5582        assert_eq!(result.close_events[0].reason, CloseReason::Manual);
5583        assert!(result.open_position_snapshots.is_empty());
5584        assert!(result.pending_order_snapshots.is_empty());
5585    }
5586
5587    #[test]
5588    fn exact_time_close_cannot_beat_later_symbol_stoploss() {
5589        let open_ts = ts(10, 0, 0);
5590        let execution_ts = ts(10, 0, 1);
5591        let events = vec![
5592            FeedEvent::new(
5593                tick("XAUUSD", 100.0, 100.0, open_ts),
5594                EventMetadata::new(SeriesRoles::PRIMARY, 1, 0),
5595            ),
5596            FeedEvent::new(
5597                tick("EURUSD", 1.1, 1.1, execution_ts),
5598                EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
5599            ),
5600            FeedEvent::new(
5601                tick("XAUUSD", 98.0, 98.0, execution_ts),
5602                EventMetadata::new(SeriesRoles::PRIMARY, 1, 1),
5603            ),
5604        ];
5605        let signals = vec![
5606            RawSignal::Entry {
5607                ts: open_ts,
5608                symbol: "XAUUSD".into(),
5609                side: Side::Buy,
5610                order_type: OrderType::Market,
5611                price: None,
5612                risk_multiplier: 1.0,
5613                stoploss: Some(99.0),
5614                targets: Vec::new(),
5615                group: None,
5616                trade_id: Some("later-stop".into()),
5617                entry_class: None,
5618            },
5619            RawSignal::Close {
5620                ts: execution_ts,
5621                position: PositionRef::ByTradeId {
5622                    trade_id: "later-stop".into(),
5623                },
5624            },
5625        ];
5626        let mut feed = VecFeed::from_feed_events(events);
5627        let result = BacktestRunner::new_future(
5628            BacktestConfig {
5629                close_on_finish: false,
5630                ..fixed_lot_config()
5631            },
5632            FutureQuoteConfig::default(),
5633        )
5634        .run_raw_signals_future(&mut feed, signals, None);
5635
5636        assert_eq!(result.close_events.len(), 1);
5637        assert_eq!(result.close_events[0].reason, CloseReason::Stoploss);
5638        assert_eq!(
5639            result.recorded_fills.last().unwrap().fill.purpose,
5640            FillPurpose::StopLoss
5641        );
5642        assert!(!result.action_dispositions.iter().any(|disposition| {
5643            disposition.action_id.starts_with("signal:00000001")
5644                && disposition.status == crate::ledger::ActionDispositionStatus::Applied
5645        }));
5646    }
5647
5648    #[test]
5649    fn exact_time_multisymbol_closes_preserve_signal_order() {
5650        let open_ts = ts(10, 0, 0);
5651        let close_ts = ts(10, 0, 1);
5652        let mut events = Vec::new();
5653        for (timestamp, row) in [(open_ts, 0), (close_ts, 1)] {
5654            events.push(FeedEvent::new(
5655                tick("EURUSD", 1.1, 1.1, timestamp),
5656                EventMetadata::new(SeriesRoles::PRIMARY, 0, row),
5657            ));
5658            events.push(FeedEvent::new(
5659                tick("XAUUSD", 100.0, 100.0, timestamp),
5660                EventMetadata::new(SeriesRoles::PRIMARY, 1, row),
5661            ));
5662        }
5663        let entry = |symbol: &str, trade_id: &str| RawSignal::Entry {
5664            ts: open_ts,
5665            symbol: symbol.into(),
5666            side: Side::Buy,
5667            order_type: OrderType::Market,
5668            price: None,
5669            risk_multiplier: 1.0,
5670            stoploss: None,
5671            targets: Vec::new(),
5672            group: None,
5673            trade_id: Some(trade_id.into()),
5674            entry_class: None,
5675        };
5676        let close = |trade_id: &str| RawSignal::Close {
5677            ts: close_ts,
5678            position: PositionRef::ByTradeId {
5679                trade_id: trade_id.into(),
5680            },
5681        };
5682        let signals = vec![
5683            entry("XAUUSD", "close-first"),
5684            entry("EURUSD", "close-second"),
5685            close("close-first"),
5686            close("close-second"),
5687        ];
5688        let mut feed = VecFeed::from_feed_events(events);
5689        let result = BacktestRunner::new_future(
5690            BacktestConfig {
5691                close_on_finish: false,
5692                ..fixed_lot_config()
5693            },
5694            FutureQuoteConfig::default(),
5695        )
5696        .run_raw_signals_future(&mut feed, signals, None);
5697
5698        assert_eq!(
5699            result
5700                .close_events
5701                .iter()
5702                .map(|event| event.symbol.as_str())
5703                .collect::<Vec<_>>(),
5704            vec!["XAUUSD", "EURUSD"]
5705        );
5706        assert!(
5707            result
5708                .close_events
5709                .iter()
5710                .all(|event| event.reason == CloseReason::Manual)
5711        );
5712    }
5713
5714    #[test]
5715    fn future_streaming_quiescence_waits_for_all_blockers() {
5716        let run = |events: Vec<MarketEvent>, signals: Vec<RawSignal>, config: BacktestConfig| {
5717            let polls = std::rc::Rc::new(std::cell::Cell::new(0));
5718            let primary_eod = events.last().map(MarketEvent::ts);
5719            let mut feed = CountingBatchFeed {
5720                batches: events.into_iter().map(primary_batch).collect(),
5721                polls: polls.clone(),
5722            };
5723            BacktestRunner::new_future(config, FutureQuoteConfig::default())
5724                .run_raw_signals_future_streaming_controlled(
5725                    &mut feed,
5726                    primary_eod,
5727                    signals,
5728                    None,
5729                    || false,
5730                    |_| {},
5731                )
5732                .unwrap();
5733            polls.get()
5734        };
5735        let eur_events = vec![
5736            tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5737            tick("EURUSD", 1.1001, 1.1003, ts(10, 0, 1)),
5738            tick("EURUSD", 1.1002, 1.1004, ts(10, 0, 2)),
5739        ];
5740
5741        let immediately_quiescent = run(
5742            eur_events.clone(),
5743            vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
5744            BacktestConfig::default(),
5745        );
5746        assert_eq!(immediately_quiescent, 1);
5747
5748        let scheduled = run(
5749            eur_events.clone(),
5750            vec![RawSignal::CloseAll { ts: ts(10, 0, 2) }],
5751            BacktestConfig::default(),
5752        );
5753        assert_eq!(scheduled, 3, "scheduled signals must block termination");
5754
5755        let mut two_symbol_config = BacktestConfig {
5756            close_on_finish: false,
5757            ..fixed_lot_config()
5758        };
5759        two_symbol_config
5760            .symbol_specs
5761            .insert("GBPUSD".into(), test_symbol_spec("GBPUSD"));
5762        let queued = run(
5763            vec![
5764                tick("EURUSD", 1.1000, 1.1002, ts(10, 0, 0)),
5765                tick("GBPUSD", 1.2500, 1.2502, ts(10, 0, 1)),
5766                tick("GBPUSD", 1.2501, 1.2503, ts(10, 0, 2)),
5767            ],
5768            vec![
5769                market_entry(ts(10, 0, 0), "GBPUSD", OrderType::Market),
5770                RawSignal::CloseAll { ts: ts(10, 0, 1) },
5771            ],
5772            two_symbol_config,
5773        );
5774        assert_eq!(
5775            queued, 2,
5776            "queued actions must wait for an eligible symbol quote"
5777        );
5778
5779        let open = run(
5780            eur_events.clone(),
5781            vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Market)],
5782            BacktestConfig {
5783                close_on_finish: false,
5784                ..fixed_lot_config()
5785            },
5786        );
5787        assert_eq!(
5788            open, 4,
5789            "open positions must consume the stream through EOD"
5790        );
5791
5792        let pending = run(
5793            eur_events,
5794            vec![market_entry(ts(10, 0, 0), "EURUSD", OrderType::Limit)],
5795            fixed_lot_config(),
5796        );
5797        assert_eq!(
5798            pending, 4,
5799            "pending orders must consume the stream through EOD"
5800        );
5801    }
5802
5803    #[test]
5804    fn future_mtm_output_policies_bound_curve_and_validate_before_feed_use() {
5805        assert_eq!(
5806            FutureQuoteConfig::default().mtm_output,
5807            MtmOutputPolicy::Bounded { max_points: 4_096 }
5808        );
5809        let events: Vec<_> = (0..12)
5810            .map(|second| tick("EURUSD", 100.0, 100.0, ts(10, 0, second)))
5811            .collect();
5812
5813        let pending = RawSignal::Entry {
5814            ts: ts(10, 0, 0),
5815            symbol: "EURUSD".into(),
5816            side: Side::Buy,
5817            order_type: OrderType::Limit,
5818            price: Some(90.0),
5819            risk_multiplier: 1.0,
5820            stoploss: None,
5821            targets: Vec::new(),
5822            group: None,
5823            trade_id: Some("mtm-policy-blocker".into()),
5824            entry_class: None,
5825        };
5826        let run = |policy| {
5827            let mut feed = VecFeed::new(events.clone());
5828            BacktestRunner::new_future(
5829                fixed_lot_config(),
5830                FutureQuoteConfig {
5831                    mtm_output: policy,
5832                    ..FutureQuoteConfig::default()
5833                },
5834            )
5835            .run_raw_signals_future(&mut feed, vec![pending.clone()], None)
5836        };
5837
5838        let none = run(MtmOutputPolicy::None);
5839        assert!(none.mtm_equity_curve.is_empty());
5840        assert_eq!(none.mtm_output_summary.observed_points, 13);
5841        assert_eq!(none.mtm_output_summary.omitted_points, 13);
5842
5843        let bounded = run(MtmOutputPolicy::Bounded { max_points: 8 });
5844        assert_eq!(bounded.mtm_equity_curve.len(), 8);
5845        assert_eq!(bounded.mtm_output_summary.observed_points, 13);
5846        assert_eq!(bounded.mtm_output_summary.retained_points, 8);
5847        assert_eq!(bounded.mtm_output_summary.omitted_points, 5);
5848
5849        let full = run(MtmOutputPolicy::Full);
5850        assert_eq!(full.mtm_equity_curve.len(), 13);
5851        assert_eq!(full.mtm_output_summary.observed_points, 13);
5852        assert_eq!(full.mtm_output_summary.omitted_points, 0);
5853        assert_eq!(
5854            full.mtm_equity_curve
5855                .iter()
5856                .filter(|point| {
5857                    point.observation_kind.as_deref()
5858                        == Some(EquityObservationKind::PostOutput.as_str())
5859                })
5860                .count(),
5861            0
5862        );
5863
5864        let mut invalid_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5865        let rejected = BacktestRunner::new_future(
5866            BacktestConfig::default(),
5867            FutureQuoteConfig {
5868                mtm_output: MtmOutputPolicy::Bounded { max_points: 7 },
5869                ..FutureQuoteConfig::default()
5870            },
5871        )
5872        .run_raw_signals_future(&mut invalid_feed, Vec::new(), None);
5873        assert_eq!(invalid_feed.remaining(), 1);
5874        assert!(rejected.action_dispositions.iter().any(|disposition| {
5875            disposition.action_id == "configuration"
5876                && disposition
5877                    .reason
5878                    .as_deref()
5879                    .is_some_and(|reason| reason.contains("MTM max_points"))
5880        }));
5881    }
5882
5883    #[test]
5884    fn future_mtm_records_changed_post_output_observation_kind() {
5885        let mut feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
5886        let signal = RawSignal::Entry {
5887            ts: ts(10, 0, 0),
5888            symbol: "EURUSD".into(),
5889            side: Side::Buy,
5890            order_type: OrderType::Market,
5891            price: None,
5892            risk_multiplier: 1.0,
5893            stoploss: None,
5894            targets: Vec::new(),
5895            group: None,
5896            trade_id: Some("mtm-kind".into()),
5897            entry_class: None,
5898        };
5899        let result = BacktestRunner::new_future(
5900            BacktestConfig {
5901                close_on_finish: false,
5902                ..fixed_lot_config()
5903            },
5904            FutureQuoteConfig {
5905                mtm_output: MtmOutputPolicy::Full,
5906                ..FutureQuoteConfig::default()
5907            },
5908        )
5909        .run_raw_signals_future(&mut feed, vec![signal], None);
5910
5911        let kinds: Vec<_> = result
5912            .mtm_equity_curve
5913            .iter()
5914            .filter_map(|point| point.observation_kind.as_deref())
5915            .collect();
5916        assert_eq!(
5917            kinds,
5918            vec![
5919                EquityObservationKind::PreSettlement.as_str(),
5920                EquityObservationKind::PostOutput.as_str(),
5921                EquityObservationKind::EndOfData.as_str(),
5922            ]
5923        );
5924        assert_eq!(
5925            result
5926                .execution_metadata
5927                .as_ref()
5928                .unwrap()
5929                .tags
5930                .get("termination_reason")
5931                .map(String::as_str),
5932            Some("end_of_data")
5933        );
5934    }
5935
5936    #[test]
5937    fn future_fallible_batch_feed_propagates_source_error() {
5938        let batch = TimestampBatch {
5939            ts: ts(10, 0, 0),
5940            events: vec![FeedEvent::new(
5941                tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
5942                EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
5943            )],
5944        };
5945        let mut feed = ScriptedBatchFeed {
5946            batches: VecDeque::from([Ok(Some(batch)), Err("feed failed")]),
5947        };
5948        let result =
5949            BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
5950                .run_raw_signals_future_fallible(&mut feed, Vec::new(), None);
5951
5952        assert!(matches!(result, Err("feed failed")));
5953    }
5954
5955    // ── Simple strategy for testing ─────────────────────────────────────
5956
5957    /// Buys on the first tick, with SL and TP.
5958    struct BuyOnceStrategy {
5959        entered: bool,
5960    }
5961
5962    impl BuyOnceStrategy {
5963        fn new() -> Self {
5964            Self { entered: false }
5965        }
5966    }
5967
5968    impl Strategy for BuyOnceStrategy {
5969        fn on_event(&mut self, event: &MarketEvent) -> Vec<Action> {
5970            if self.entered {
5971                return vec![];
5972            }
5973            if let MarketEvent::Tick { symbol, ask, .. } = event {
5974                self.entered = true;
5975                vec![Action::Open {
5976                    symbol: symbol.clone(),
5977                    side: Side::Buy,
5978                    order_type: OrderType::Market,
5979                    price: Some(*ask),
5980                    size: 1.0,
5981                    stoploss: Some(*ask - 0.0050),
5982                    targets: vec![TargetSpec {
5983                        price: *ask + 0.0050,
5984                        close_ratio: 1.0,
5985                    }],
5986                    rules: vec![],
5987                    group: None,
5988                    trade_id: None,
5989                }]
5990            } else {
5991                vec![]
5992            }
5993        }
5994
5995        fn on_finished(&mut self) -> Vec<Action> {
5996            // Don't close — let close_on_finish handle it if TP/SL haven't
5997            // triggered.
5998            vec![]
5999        }
6000    }
6001
6002    // ── Strategy-driven tests ───────────────────────────────────────────
6003
6004    #[test]
6005    fn strategy_backtest_tp_hit() {
6006        let events = vec![
6007            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6008            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6009            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6010            tick("EURUSD", 1.0890, 1.0892, ts(10, 0, 3)),
6011            // TP at 1.0900 (entry 1.0850 + 0.005)
6012            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 4)),
6013        ];
6014        let mut feed = VecFeed::new(events);
6015        let mut strategy = BuyOnceStrategy::new();
6016
6017        let config = BacktestConfig {
6018            initial_balance: 10_000.0,
6019            close_on_finish: true,
6020            ..Default::default()
6021        };
6022        let runner = BacktestRunner::new(config);
6023        let result = runner.run_strategy(&mut feed, &mut strategy);
6024
6025        assert_eq!(result.total_trades, 1);
6026        assert_eq!(result.winning_trades, 1);
6027        assert!(result.total_pnl > 0.0);
6028        assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6029    }
6030
6031    #[test]
6032    fn strategy_backtest_sl_hit() {
6033        let events = vec![
6034            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6035            tick("EURUSD", 1.0830, 1.0832, ts(10, 0, 1)),
6036            // SL at 1.0800 (entry 1.0850 - 0.005)
6037            tick("EURUSD", 1.0799, 1.0801, ts(10, 0, 2)),
6038        ];
6039        let mut feed = VecFeed::new(events);
6040        let mut strategy = BuyOnceStrategy::new();
6041
6042        let runner = BacktestRunner::with_defaults();
6043        let result = runner.run_strategy(&mut feed, &mut strategy);
6044
6045        assert_eq!(result.total_trades, 1);
6046        assert_eq!(result.losing_trades, 1);
6047        assert!(result.total_pnl < 0.0);
6048        assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6049    }
6050
6051    #[test]
6052    fn strategy_close_on_finish() {
6053        // Price never reaches TP or SL — position should be closed at end.
6054        let events = vec![
6055            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6056            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6057            tick("EURUSD", 1.0852, 1.0854, ts(10, 0, 2)),
6058        ];
6059        let mut feed = VecFeed::new(events);
6060        let mut strategy = BuyOnceStrategy::new();
6061
6062        let config = BacktestConfig {
6063            initial_balance: 10_000.0,
6064            close_on_finish: true,
6065            ..Default::default()
6066        };
6067        let runner = BacktestRunner::new(config);
6068        let result = runner.run_strategy(&mut feed, &mut strategy);
6069
6070        assert_eq!(result.total_trades, 1);
6071        assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6072    }
6073
6074    #[test]
6075    fn strategy_no_close_on_finish() {
6076        let events = vec![
6077            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6078            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6079        ];
6080        let mut feed = VecFeed::new(events);
6081        let mut strategy = BuyOnceStrategy::new();
6082
6083        let config = BacktestConfig {
6084            initial_balance: 10_000.0,
6085            close_on_finish: false,
6086            ..Default::default()
6087        };
6088        let runner = BacktestRunner::new(config);
6089        let result = runner.run_strategy(&mut feed, &mut strategy);
6090
6091        // Position left open — no trades recorded.
6092        assert_eq!(result.total_trades, 0);
6093    }
6094
6095    // ── Raw signal replay tests ─────────────────────────────────────────
6096
6097    #[test]
6098    fn legacy_unprofiled_targets_default_to_equal_weights() {
6099        let events = vec![
6100            tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
6101            tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
6102            tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
6103        ];
6104        let mut feed = VecFeed::new(events);
6105        let signals = vec![RawSignal::Entry {
6106            ts: ts(10, 0, 0),
6107            symbol: "EURUSD".into(),
6108            side: Side::Buy,
6109            order_type: OrderType::Market,
6110            price: Some(1.0000),
6111            risk_multiplier: 1.0,
6112            stoploss: None,
6113            targets: vec![1.1000, 1.2000],
6114            group: None,
6115            trade_id: Some("equal-targets".into()),
6116            entry_class: None,
6117        }];
6118
6119        let result = BacktestRunner::new(BacktestConfig {
6120            close_on_finish: false,
6121            ..fixed_lot_config()
6122        })
6123        .run_raw_signals(&mut feed, signals, None);
6124
6125        assert_eq!(result.trade_log.len(), 2);
6126        assert!(
6127            result
6128                .trade_log
6129                .iter()
6130                .all(|trade| (trade.size - 0.5).abs() < f64::EPSILON)
6131        );
6132        assert!(
6133            result
6134                .trade_log
6135                .iter()
6136                .all(|trade| trade.close_reason == CloseReason::Target)
6137        );
6138    }
6139
6140    #[test]
6141    fn legacy_atomic_target_modification_retains_profile_ratio() {
6142        let events = vec![
6143            tick("EURUSD", 1.0000, 1.0000, ts(10, 0, 0)),
6144            tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
6145            tick("EURUSD", 1.2000, 1.2000, ts(10, 0, 2)),
6146            tick("EURUSD", 1.3000, 1.3000, ts(10, 0, 3)),
6147        ];
6148        let mut feed = VecFeed::new(events);
6149        let position = PositionRef::ByTradeId {
6150            trade_id: "modified-target".into(),
6151        };
6152        let signals = vec![
6153            RawSignal::Entry {
6154                ts: ts(10, 0, 0),
6155                symbol: "EURUSD".into(),
6156                side: Side::Buy,
6157                order_type: OrderType::Market,
6158                price: Some(1.0000),
6159                risk_multiplier: 1.0,
6160                stoploss: None,
6161                targets: vec![1.1000, 1.3000],
6162                group: None,
6163                trade_id: Some("modified-target".into()),
6164                entry_class: None,
6165            },
6166            RawSignal::ModifyTarget {
6167                ts: ts(10, 0, 0),
6168                position,
6169                old_price: 1.1000,
6170                new_price: 1.2000,
6171            },
6172        ];
6173        let profile = ManagementProfile {
6174            name: "non-default-ratios".into(),
6175            target_selection: None,
6176            use_targets: vec![1, 2],
6177            close_ratios: vec![0.25, 0.75],
6178            target_source: TargetSource::FromSignal,
6179            stoploss_mode: StoplossMode::FromSignal,
6180            rules: vec![],
6181            group_override: None,
6182            let_remainder_run: false,
6183            entry_geometry: EntryGeometryPolicy::Strict,
6184        };
6185
6186        let result = BacktestRunner::new(BacktestConfig {
6187            close_on_finish: false,
6188            ..fixed_lot_config()
6189        })
6190        .run_raw_signals(&mut feed, signals, Some(&profile));
6191
6192        assert_eq!(result.trade_log.len(), 2);
6193        assert!((result.trade_log[0].exit_price - 1.2000).abs() < f64::EPSILON);
6194        assert!((result.trade_log[0].size - 0.25).abs() < f64::EPSILON);
6195        assert!((result.trade_log[1].exit_price - 1.3000).abs() < f64::EPSILON);
6196        assert!((result.trade_log[1].size - 0.75).abs() < f64::EPSILON);
6197    }
6198
6199    #[test]
6200    fn run_raw_signals_entry_only() {
6201        let events = vec![
6202            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6203            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6204            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6205            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
6206        ];
6207        let mut feed = VecFeed::new(events);
6208
6209        let raw_signals = vec![RawSignal::Entry {
6210            ts: ts(10, 0, 0),
6211            symbol: "EURUSD".into(),
6212            side: Side::Buy,
6213            order_type: OrderType::Market,
6214            price: Some(1.0850),
6215            risk_multiplier: 1.0,
6216            stoploss: Some(1.0800),
6217            targets: vec![1.0900],
6218            group: None,
6219            trade_id: None,
6220            entry_class: None,
6221        }];
6222
6223        let runner = BacktestRunner::new(fixed_lot_config());
6224        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6225
6226        assert_eq!(result.total_trades, 1);
6227        assert_eq!(result.winning_trades, 1);
6228    }
6229
6230    #[test]
6231    fn run_raw_signals_open_then_close() {
6232        let events = vec![
6233            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6234            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6235            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6236            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6237        ];
6238        let mut feed = VecFeed::new(events);
6239
6240        let raw_signals = vec![
6241            RawSignal::Entry {
6242                ts: ts(10, 0, 0),
6243                symbol: "EURUSD".into(),
6244                side: Side::Buy,
6245                order_type: OrderType::Market,
6246                price: Some(1.0850),
6247                risk_multiplier: 1.0,
6248                stoploss: None,
6249                targets: vec![],
6250                group: None,
6251                trade_id: Some("t1".into()),
6252                entry_class: None,
6253            },
6254            RawSignal::Close {
6255                ts: ts(10, 0, 2),
6256                position: PositionRef::ByTradeId {
6257                    trade_id: "t1".into(),
6258                },
6259            },
6260        ];
6261
6262        let config = BacktestConfig {
6263            initial_balance: 10_000.0,
6264            close_on_finish: false,
6265            ..fixed_lot_config()
6266        };
6267        let runner = BacktestRunner::new(config);
6268        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6269
6270        assert_eq!(result.total_trades, 1);
6271        assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6272    }
6273
6274    #[test]
6275    fn run_raw_signals_open_then_modify_sl() {
6276        // Open a position, then move SL closer. If price drops to new SL, it triggers.
6277        let events = vec![
6278            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6279            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
6280            // SL modify happens at ts(10,0,2)
6281            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 2)),
6282            // Price drops to modified SL at 1.0840
6283            tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
6284        ];
6285        let mut feed = VecFeed::new(events);
6286
6287        let raw_signals = vec![
6288            RawSignal::Entry {
6289                ts: ts(10, 0, 0),
6290                symbol: "EURUSD".into(),
6291                side: Side::Buy,
6292                order_type: OrderType::Market,
6293                price: Some(1.0850),
6294                risk_multiplier: 1.0,
6295                stoploss: Some(1.0800),
6296                targets: vec![],
6297                group: None,
6298                trade_id: Some("t1".into()),
6299                entry_class: None,
6300            },
6301            RawSignal::ModifyStoploss {
6302                ts: ts(10, 0, 2),
6303                position: PositionRef::ByTradeId {
6304                    trade_id: "t1".into(),
6305                },
6306                price: 1.0840,
6307            },
6308        ];
6309
6310        let config = BacktestConfig {
6311            initial_balance: 10_000.0,
6312            close_on_finish: true,
6313            ..fixed_lot_config()
6314        };
6315        let runner = BacktestRunner::new(config);
6316        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6317
6318        assert_eq!(result.total_trades, 1);
6319        assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6320    }
6321
6322    #[test]
6323    fn run_raw_signals_open_then_partial_close() {
6324        let events = vec![
6325            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6326            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 1)),
6327            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6328            tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 3)),
6329        ];
6330        let mut feed = VecFeed::new(events);
6331
6332        let raw_signals = vec![
6333            RawSignal::Entry {
6334                ts: ts(10, 0, 0),
6335                symbol: "EURUSD".into(),
6336                side: Side::Buy,
6337                order_type: OrderType::Market,
6338                price: Some(1.0850),
6339                risk_multiplier: 1.0,
6340                stoploss: None,
6341                targets: vec![],
6342                group: None,
6343                trade_id: Some("t1".into()),
6344                entry_class: None,
6345            },
6346            RawSignal::ClosePartial {
6347                ts: ts(10, 0, 1),
6348                position: PositionRef::ByTradeId {
6349                    trade_id: "t1".into(),
6350                },
6351                ratio: 0.5,
6352            },
6353        ];
6354
6355        let config = BacktestConfig {
6356            initial_balance: 10_000.0,
6357            close_on_finish: true,
6358            ..fixed_lot_config()
6359        };
6360        let runner = BacktestRunner::new(config);
6361        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6362
6363        // At least 1 trade closed (partial close + close_on_finish for remainder)
6364        assert!(result.total_trades >= 1);
6365    }
6366
6367    #[test]
6368    fn run_raw_signals_group_workflow() {
6369        let events = vec![
6370            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6371            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6372            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6373            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6374            tick("EURUSD", 1.0880, 1.0882, ts(10, 0, 4)),
6375        ];
6376        let mut feed = VecFeed::new(events);
6377
6378        let raw_signals = vec![
6379            // Open 2 positions in same group
6380            RawSignal::Entry {
6381                ts: ts(10, 0, 0),
6382                symbol: "EURUSD".into(),
6383                side: Side::Buy,
6384                order_type: OrderType::Market,
6385                price: Some(1.0850),
6386                risk_multiplier: 1.0,
6387                stoploss: None,
6388                targets: vec![],
6389                group: Some("grp1".into()),
6390                trade_id: Some("t1".into()),
6391                entry_class: None,
6392            },
6393            RawSignal::Entry {
6394                ts: ts(10, 0, 1),
6395                symbol: "EURUSD".into(),
6396                side: Side::Buy,
6397                order_type: OrderType::Market,
6398                price: Some(1.0857),
6399                risk_multiplier: 1.0,
6400                stoploss: None,
6401                targets: vec![],
6402                group: Some("grp1".into()),
6403                trade_id: Some("t2".into()),
6404                entry_class: None,
6405            },
6406            // Close entire group
6407            RawSignal::CloseAllInGroup {
6408                ts: ts(10, 0, 3),
6409                group_id: "grp1".into(),
6410            },
6411        ];
6412
6413        let config = BacktestConfig {
6414            initial_balance: 10_000.0,
6415            close_on_finish: false,
6416            ..fixed_lot_config()
6417        };
6418        let runner = BacktestRunner::new(config);
6419        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6420
6421        assert_eq!(result.total_trades, 2);
6422    }
6423
6424    #[test]
6425    fn run_raw_signals_close_all_of_symbol() {
6426        let events = vec![
6427            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6428            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6429            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6430            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6431        ];
6432        let mut feed = VecFeed::new(events);
6433
6434        let raw_signals = vec![
6435            RawSignal::Entry {
6436                ts: ts(10, 0, 0),
6437                symbol: "EURUSD".into(),
6438                side: Side::Buy,
6439                order_type: OrderType::Market,
6440                price: Some(1.0850),
6441                risk_multiplier: 1.0,
6442                stoploss: None,
6443                targets: vec![],
6444                group: None,
6445                trade_id: Some("t1".into()),
6446                entry_class: None,
6447            },
6448            RawSignal::Entry {
6449                ts: ts(10, 0, 0),
6450                symbol: "EURUSD".into(),
6451                side: Side::Buy,
6452                order_type: OrderType::Market,
6453                price: Some(1.0850),
6454                risk_multiplier: 0.5,
6455                stoploss: None,
6456                targets: vec![],
6457                group: None,
6458                trade_id: Some("t2".into()),
6459                entry_class: None,
6460            },
6461            RawSignal::CloseAllOf {
6462                ts: ts(10, 0, 2),
6463                symbol: "EURUSD".into(),
6464            },
6465        ];
6466
6467        let config = BacktestConfig {
6468            initial_balance: 10_000.0,
6469            close_on_finish: false,
6470            ..fixed_lot_config()
6471        };
6472        let runner = BacktestRunner::new(config);
6473        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6474
6475        assert_eq!(result.total_trades, 2);
6476    }
6477
6478    #[test]
6479    fn run_raw_signals_with_profile() {
6480        let events = vec![
6481            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6482            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6483            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
6484        ];
6485        let mut feed = VecFeed::new(events);
6486
6487        let profile = ManagementProfile {
6488            name: "test".into(),
6489            target_selection: None,
6490            use_targets: vec![1],
6491            close_ratios: vec![1.0],
6492            target_source: TargetSource::FromSignal,
6493            stoploss_mode: StoplossMode::FromSignal,
6494            rules: vec![],
6495            group_override: None,
6496            let_remainder_run: false,
6497            entry_geometry: EntryGeometryPolicy::Strict,
6498        };
6499
6500        let raw_signals = vec![RawSignal::Entry {
6501            ts: ts(10, 0, 0),
6502            symbol: "EURUSD".into(),
6503            side: Side::Buy,
6504            order_type: OrderType::Market,
6505            price: Some(1.0850),
6506            risk_multiplier: 1.0,
6507            stoploss: Some(1.0800),
6508            targets: vec![1.0900],
6509            group: None,
6510            trade_id: Some("t1".into()),
6511            entry_class: None,
6512        }];
6513
6514        let runner = BacktestRunner::new(fixed_lot_config());
6515        let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6516
6517        assert_eq!(result.total_trades, 1);
6518        assert_eq!(result.winning_trades, 1);
6519        assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6520    }
6521
6522    #[test]
6523    fn run_raw_signals_with_profile_preserves_trade_id() {
6524        // Regression for profile-supplied trade ID propagation.
6525        // A raw entry must expose its trade ID so a later PositionRef::ByTradeId signal can resolve and close the position.
6526        let events = vec![
6527            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6528            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6529        ];
6530        let mut feed = VecFeed::new(events);
6531
6532        let profile = ManagementProfile {
6533            name: "test".into(),
6534            target_selection: None,
6535            use_targets: vec![1],
6536            close_ratios: vec![1.0],
6537            target_source: TargetSource::FromSignal,
6538            stoploss_mode: StoplossMode::FromSignal,
6539            rules: vec![],
6540            group_override: None,
6541            let_remainder_run: false,
6542            entry_geometry: EntryGeometryPolicy::Strict,
6543        };
6544
6545        let raw_signals = vec![
6546            RawSignal::Entry {
6547                ts: ts(10, 0, 0),
6548                symbol: "EURUSD".into(),
6549                side: Side::Buy,
6550                order_type: OrderType::Market,
6551                price: Some(1.0850),
6552                risk_multiplier: 1.0,
6553                stoploss: Some(1.0800),
6554                targets: vec![1.0900],
6555                group: None,
6556                trade_id: Some("msg-100".into()),
6557                entry_class: None,
6558            },
6559            RawSignal::Close {
6560                ts: ts(10, 0, 1),
6561                position: PositionRef::ByTradeId {
6562                    trade_id: "msg-100".into(),
6563                },
6564            },
6565        ];
6566
6567        let runner = BacktestRunner::new(fixed_lot_config());
6568        let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6569
6570        assert_eq!(result.total_trades, 1);
6571        assert_eq!(result.trade_log[0].close_reason, CloseReason::Manual);
6572    }
6573
6574    #[test]
6575    fn run_raw_signals_no_profile() {
6576        // Without a profile, entry signals are converted directly.
6577        let events = vec![
6578            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6579            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6580            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 2)),
6581        ];
6582        let mut feed = VecFeed::new(events);
6583
6584        let raw_signals = vec![RawSignal::Entry {
6585            ts: ts(10, 0, 0),
6586            symbol: "EURUSD".into(),
6587            side: Side::Buy,
6588            order_type: OrderType::Market,
6589            price: Some(1.0850),
6590            risk_multiplier: 1.0,
6591            stoploss: None,
6592            targets: vec![],
6593            group: None,
6594            trade_id: None,
6595            entry_class: None,
6596        }];
6597
6598        let config = BacktestConfig {
6599            initial_balance: 10_000.0,
6600            close_on_finish: true,
6601            ..fixed_lot_config()
6602        };
6603        let runner = BacktestRunner::new(config);
6604        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6605
6606        assert_eq!(result.total_trades, 1);
6607    }
6608
6609    #[test]
6610    fn run_raw_signals_last_on_symbol_resolution() {
6611        // Open two positions, then close the last one by symbol ref.
6612        let events = vec![
6613            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6614            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6615            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6616            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6617        ];
6618        let mut feed = VecFeed::new(events);
6619
6620        let raw_signals = vec![
6621            RawSignal::Entry {
6622                ts: ts(10, 0, 0),
6623                symbol: "EURUSD".into(),
6624                side: Side::Buy,
6625                order_type: OrderType::Market,
6626                price: Some(1.0850),
6627                risk_multiplier: 1.0,
6628                stoploss: None,
6629                targets: vec![],
6630                group: None,
6631                trade_id: Some("t1".into()),
6632                entry_class: None,
6633            },
6634            RawSignal::Entry {
6635                ts: ts(10, 0, 1),
6636                symbol: "EURUSD".into(),
6637                side: Side::Buy,
6638                order_type: OrderType::Market,
6639                price: Some(1.0857),
6640                risk_multiplier: 1.0,
6641                stoploss: None,
6642                targets: vec![],
6643                group: None,
6644                trade_id: Some("t2".into()),
6645                entry_class: None,
6646            },
6647            // Close only the second opened position via its trade_id
6648            RawSignal::Close {
6649                ts: ts(10, 0, 2),
6650                position: PositionRef::ByTradeId {
6651                    trade_id: "t2".into(),
6652                },
6653            },
6654        ];
6655
6656        let config = BacktestConfig {
6657            initial_balance: 10_000.0,
6658            close_on_finish: true,
6659            ..fixed_lot_config()
6660        };
6661        let runner = BacktestRunner::new(config);
6662        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6663
6664        // 2 trades total: one closed by signal, one by close_on_finish
6665        assert_eq!(result.total_trades, 2);
6666    }
6667
6668    #[test]
6669    fn run_raw_signals_unresolved_ref_skipped() {
6670        // Try to close a position that doesn't exist — should be silently skipped.
6671        let events = vec![
6672            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6673            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6674        ];
6675        let mut feed = VecFeed::new(events);
6676
6677        let raw_signals = vec![RawSignal::Close {
6678            ts: ts(10, 0, 0),
6679            position: PositionRef::ByTradeId {
6680                trade_id: "nonexistent".into(),
6681            },
6682        }];
6683
6684        let config = BacktestConfig {
6685            initial_balance: 10_000.0,
6686            close_on_finish: false,
6687            ..Default::default()
6688        };
6689        let runner = BacktestRunner::new(config);
6690        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6691
6692        // No positions were opened or closed.
6693        assert_eq!(result.total_trades, 0);
6694    }
6695
6696    // ── Signal replay tests ─────────────────────────────────────────────
6697
6698    #[test]
6699    fn signal_replay_basic() {
6700        let events = vec![
6701            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6702            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6703            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6704            // TP at 1.0900
6705            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 3)),
6706        ];
6707        let mut feed = VecFeed::new(events);
6708
6709        let raw_signals = vec![RawSignal::Entry {
6710            ts: ts(10, 0, 0),
6711            symbol: "EURUSD".into(),
6712            side: Side::Buy,
6713            order_type: OrderType::Market,
6714            price: Some(1.0850),
6715            risk_multiplier: 1.0,
6716            stoploss: Some(1.0800),
6717            targets: vec![1.0900],
6718            group: None,
6719            trade_id: None,
6720            entry_class: None,
6721        }];
6722
6723        let runner = BacktestRunner::new(fixed_lot_config());
6724        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6725
6726        assert_eq!(result.total_trades, 1);
6727        assert_eq!(result.winning_trades, 1);
6728        assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6729    }
6730
6731    #[test]
6732    fn signal_replay_multiple_signals() {
6733        let events = vec![
6734            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6735            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6736            // TP1 hit for first position
6737            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 2)),
6738            tick("EURUSD", 1.0910, 1.0912, ts(10, 0, 3)),
6739            tick("EURUSD", 1.0920, 1.0922, ts(10, 0, 4)),
6740        ];
6741        let mut feed = VecFeed::new(events);
6742
6743        let raw_signals = vec![
6744            RawSignal::Entry {
6745                ts: ts(10, 0, 0),
6746                symbol: "EURUSD".into(),
6747                side: Side::Buy,
6748                order_type: OrderType::Market,
6749                price: Some(1.0850),
6750                risk_multiplier: 1.0,
6751                stoploss: Some(1.0800),
6752                targets: vec![1.0900],
6753                group: None,
6754                trade_id: Some("t1".into()),
6755                entry_class: None,
6756            },
6757            RawSignal::Entry {
6758                ts: ts(10, 0, 1),
6759                symbol: "EURUSD".into(),
6760                side: Side::Buy,
6761                order_type: OrderType::Market,
6762                price: Some(1.0857),
6763                risk_multiplier: 1.0,
6764                stoploss: None,
6765                targets: vec![],
6766                group: None,
6767                trade_id: Some("t2".into()),
6768                entry_class: None,
6769            },
6770        ];
6771
6772        let config = BacktestConfig {
6773            initial_balance: 10_000.0,
6774            close_on_finish: true,
6775            ..fixed_lot_config()
6776        };
6777        let runner = BacktestRunner::new(config);
6778        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6779
6780        // First position closed by TP, second by close_on_finish
6781        assert!(result.total_trades >= 2);
6782    }
6783
6784    #[test]
6785    fn signal_replay_signal_before_data_filtered() {
6786        // Signal timestamp is before first data event.
6787        // The runner itself does not filter; the server is responsible
6788        // for date filtering. This test verifies that when a pre-window
6789        // signal IS passed to the runner, it is injected at the first
6790        // event (backward-compatible library behavior).
6791        // Server-side filtering is tested separately.
6792        let events = vec![
6793            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6794            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6795        ];
6796        let mut feed = VecFeed::new(events);
6797
6798        let raw_signals = vec![RawSignal::Entry {
6799            ts: ts(9, 0, 0), // before first tick
6800            symbol: "EURUSD".into(),
6801            side: Side::Buy,
6802            order_type: OrderType::Market,
6803            price: Some(1.0850),
6804            risk_multiplier: 1.0,
6805            stoploss: None,
6806            targets: vec![1.0900],
6807            group: None,
6808            trade_id: None,
6809            entry_class: None,
6810        }];
6811
6812        let runner = BacktestRunner::new(fixed_lot_config());
6813        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
6814
6815        assert_eq!(result.total_trades, 1);
6816        assert_eq!(result.trade_log[0].close_reason, CloseReason::Target);
6817    }
6818
6819    #[test]
6820    fn empty_feed_empty_result() {
6821        let mut feed = VecFeed::new(vec![]);
6822        let mut strategy = BuyOnceStrategy::new();
6823
6824        let runner = BacktestRunner::with_defaults();
6825        let result = runner.run_strategy(&mut feed, &mut strategy);
6826
6827        assert_eq!(result.total_trades, 0);
6828        assert!((result.final_balance - 10_000.0).abs() < f64::EPSILON);
6829    }
6830
6831    #[test]
6832    fn report_display_does_not_panic() {
6833        let events = vec![
6834            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6835            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
6836        ];
6837        let mut feed = VecFeed::new(events);
6838        let mut strategy = BuyOnceStrategy::new();
6839
6840        let runner = BacktestRunner::with_defaults();
6841        let result = runner.run_strategy(&mut feed, &mut strategy);
6842
6843        let _display = format!("{result}");
6844    }
6845
6846    #[test]
6847    fn run_raw_signals_with_profile_open_then_modify_sl_by_trade_id() {
6848        let events = vec![
6849            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6850            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6851            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6852            tick("EURUSD", 1.0838, 1.0840, ts(10, 0, 3)),
6853        ];
6854        let mut feed = VecFeed::new(events);
6855
6856        let profile = ManagementProfile {
6857            name: "test".into(),
6858            target_selection: None,
6859            use_targets: vec![1],
6860            close_ratios: vec![1.0],
6861            target_source: TargetSource::FromSignal,
6862            stoploss_mode: StoplossMode::FromSignal,
6863            rules: vec![],
6864            group_override: None,
6865            let_remainder_run: false,
6866            entry_geometry: EntryGeometryPolicy::Strict,
6867        };
6868
6869        let raw_signals = vec![
6870            RawSignal::Entry {
6871                ts: ts(10, 0, 0),
6872                symbol: "EURUSD".into(),
6873                side: Side::Buy,
6874                order_type: OrderType::Market,
6875                price: Some(1.0850),
6876                risk_multiplier: 1.0,
6877                stoploss: Some(1.0800),
6878                targets: vec![1.0900],
6879                group: None,
6880                trade_id: Some("t1".into()),
6881                entry_class: None,
6882            },
6883            RawSignal::ModifyStoploss {
6884                ts: ts(10, 0, 2),
6885                position: PositionRef::ByTradeId {
6886                    trade_id: "t1".into(),
6887                },
6888                price: 1.0840,
6889            },
6890        ];
6891
6892        let config = BacktestConfig {
6893            initial_balance: 10_000.0,
6894            close_on_finish: false,
6895            ..fixed_lot_config()
6896        };
6897        let runner = BacktestRunner::new(config);
6898        let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6899
6900        assert_eq!(result.total_trades, 1);
6901        assert_eq!(result.trade_log[0].close_reason, CloseReason::Stoploss);
6902    }
6903
6904    #[test]
6905    fn run_raw_signals_with_profile_open_then_close_partial_by_trade_id() {
6906        let events = vec![
6907            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6908            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6909            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6910        ];
6911        let mut feed = VecFeed::new(events);
6912
6913        let profile = ManagementProfile {
6914            name: "test".into(),
6915            target_selection: None,
6916            use_targets: vec![1],
6917            close_ratios: vec![1.0],
6918            target_source: TargetSource::FromSignal,
6919            stoploss_mode: StoplossMode::FromSignal,
6920            rules: vec![],
6921            group_override: None,
6922            let_remainder_run: false,
6923            entry_geometry: EntryGeometryPolicy::Strict,
6924        };
6925
6926        let raw_signals = vec![
6927            RawSignal::Entry {
6928                ts: ts(10, 0, 0),
6929                symbol: "EURUSD".into(),
6930                side: Side::Buy,
6931                order_type: OrderType::Market,
6932                price: Some(1.0850),
6933                risk_multiplier: 1.0,
6934                stoploss: Some(1.0800),
6935                targets: vec![1.0900],
6936                group: None,
6937                trade_id: Some("t1".into()),
6938                entry_class: None,
6939            },
6940            RawSignal::ClosePartial {
6941                ts: ts(10, 0, 1),
6942                position: PositionRef::ByTradeId {
6943                    trade_id: "t1".into(),
6944                },
6945                ratio: 0.5,
6946            },
6947        ];
6948
6949        let config = BacktestConfig {
6950            initial_balance: 10_000.0,
6951            close_on_finish: false,
6952            ..fixed_lot_config()
6953        };
6954        let runner = BacktestRunner::new(config);
6955        let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
6956
6957        // Partial close creates at least one trade.
6958        assert!(result.total_trades >= 1);
6959    }
6960
6961    #[test]
6962    fn run_raw_signals_multi_position_by_trade_id_with_profile() {
6963        // Two entries on EURUSD group "alpha" with different trade_ids.
6964        // Close ByTradeId for "t1" only. Verify only t1 closes by signal
6965        // and t2 remains to be closed by close_on_finish.
6966        let events = vec![
6967            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
6968            tick("EURUSD", 1.0855, 1.0857, ts(10, 0, 1)),
6969            tick("EURUSD", 1.0860, 1.0862, ts(10, 0, 2)),
6970            tick("EURUSD", 1.0870, 1.0872, ts(10, 0, 3)),
6971        ];
6972        let mut feed = VecFeed::new(events);
6973
6974        let profile = ManagementProfile {
6975            name: "test".into(),
6976            target_selection: None,
6977            use_targets: vec![1],
6978            close_ratios: vec![1.0],
6979            target_source: TargetSource::FromSignal,
6980            stoploss_mode: StoplossMode::FromSignal,
6981            rules: vec![],
6982            group_override: Some("alpha".into()),
6983            let_remainder_run: false,
6984            entry_geometry: EntryGeometryPolicy::Strict,
6985        };
6986
6987        let raw_signals = vec![
6988            RawSignal::Entry {
6989                ts: ts(10, 0, 0),
6990                symbol: "EURUSD".into(),
6991                side: Side::Buy,
6992                order_type: OrderType::Market,
6993                price: Some(1.0850),
6994                risk_multiplier: 1.0,
6995                stoploss: None,
6996                targets: vec![1.0910],
6997                group: None,
6998                trade_id: Some("t1".into()),
6999                entry_class: None,
7000            },
7001            RawSignal::Entry {
7002                ts: ts(10, 0, 1),
7003                symbol: "EURUSD".into(),
7004                side: Side::Buy,
7005                order_type: OrderType::Market,
7006                price: Some(1.0857),
7007                risk_multiplier: 1.0,
7008                stoploss: None,
7009                targets: vec![1.0910],
7010                group: None,
7011                trade_id: Some("t2".into()),
7012                entry_class: None,
7013            },
7014            // Close only t1 by trade_id.
7015            RawSignal::Close {
7016                ts: ts(10, 0, 2),
7017                position: PositionRef::ByTradeId {
7018                    trade_id: "t1".into(),
7019                },
7020            },
7021        ];
7022
7023        let config = BacktestConfig {
7024            initial_balance: 10_000.0,
7025            close_on_finish: true,
7026            ..fixed_lot_config()
7027        };
7028        let runner = BacktestRunner::new(config);
7029        let result = runner.run_raw_signals(&mut feed, raw_signals, Some(&profile));
7030
7031        // t1 closed by signal, t2 closed by close_on_finish = 2 total.
7032        assert_eq!(result.total_trades, 2);
7033        // Both should be in group "alpha" from profile override.
7034        for trade in &result.trade_log {
7035            assert_eq!(trade.group.as_deref(), Some("alpha"));
7036        }
7037    }
7038
7039    #[test]
7040    fn merged_feed_manual_close_uses_correct_symbol_quote() {
7041        // Regression test for Issue 1 Part 3:
7042        // Open XAUUSD, then close it manually while the current merged-feed
7043        // event is a GBPJPY tick. The exit price must be a XAUUSD price,
7044        // not a GBPJPY price.
7045        use crate::data_feed::MarketEvent;
7046        let events = vec![
7047            MarketEvent::Tick {
7048                symbol: "XAUUSD".into(),
7049                ts: ts(10, 0, 0),
7050                bid: 5000.0,
7051                ask: 5001.0,
7052            },
7053            MarketEvent::Tick {
7054                symbol: "GBPJPY".into(),
7055                ts: ts(10, 0, 1),
7056                bid: 210.0,
7057                ask: 211.0,
7058            },
7059            MarketEvent::Tick {
7060                symbol: "XAUUSD".into(),
7061                ts: ts(10, 0, 2),
7062                bid: 5050.0,
7063                ask: 5051.0,
7064            },
7065            // GBPJPY event at ts(10,0,3) - manual close fires here.
7066            MarketEvent::Tick {
7067                symbol: "GBPJPY".into(),
7068                ts: ts(10, 0, 3),
7069                bid: 212.0,
7070                ask: 213.0,
7071            },
7072        ];
7073        let mut feed = VecFeed::new(events);
7074
7075        let raw_signals = vec![
7076            RawSignal::Entry {
7077                ts: ts(10, 0, 0),
7078                symbol: "XAUUSD".into(),
7079                side: Side::Buy,
7080                order_type: OrderType::Market,
7081                price: Some(5000.0),
7082                risk_multiplier: 1.0,
7083                stoploss: None,
7084                targets: vec![],
7085                group: None,
7086                trade_id: Some("xau-1".into()),
7087                entry_class: None,
7088            },
7089            // Manual close at ts(10,0,3) while current event is GBPJPY.
7090            RawSignal::Close {
7091                ts: ts(10, 0, 3),
7092                position: PositionRef::ByTradeId {
7093                    trade_id: "xau-1".into(),
7094                },
7095            },
7096        ];
7097
7098        let config = BacktestConfig {
7099            initial_balance: 10_000.0,
7100            close_on_finish: false,
7101            ..fixed_lot_config()
7102        };
7103        let runner = BacktestRunner::new(config);
7104        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
7105
7106        assert_eq!(result.total_trades, 1);
7107        let trade = &result.trade_log[0];
7108        assert_eq!(trade.symbol, "XAUUSD");
7109        // Exit price must be a XAUUSD price (~5050), not GBPJPY (~212).
7110        assert!(
7111            trade.exit_price > 4000.0,
7112            "Exit price should be XAUUSD (~5050), got {}",
7113            trade.exit_price
7114        );
7115    }
7116
7117    fn long_tick_feed(count: usize) -> VecFeed {
7118        let start = ts(10, 0, 0);
7119        VecFeed::new(
7120            (0..count)
7121                .map(|index| {
7122                    tick(
7123                        "EURUSD",
7124                        1.0848,
7125                        1.0850,
7126                        start + Duration::milliseconds(index as i64),
7127                    )
7128                })
7129                .collect(),
7130        )
7131    }
7132
7133    #[test]
7134    fn legacy_replay_can_be_cancelled_during_event_processing() {
7135        let cancelled = std::cell::Cell::new(false);
7136        let mut feed = long_tick_feed(1_000);
7137        let outcome = BacktestRunner::with_defaults().run_raw_signals_controlled(
7138            &mut feed,
7139            Vec::new(),
7140            None,
7141            || cancelled.get(),
7142            |progress| {
7143                if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
7144                    cancelled.set(true);
7145                }
7146            },
7147        );
7148
7149        assert_eq!(outcome.unwrap_err(), ReplayCancelled);
7150        assert!(
7151            feed.remaining() > 0,
7152            "cancellation must stop further replay"
7153        );
7154    }
7155
7156    #[test]
7157    fn future_quote_replay_can_be_cancelled_during_event_processing() {
7158        let cancelled = std::cell::Cell::new(false);
7159        let mut feed = long_tick_feed(1_000);
7160        let runner = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default());
7161        let pending = RawSignal::Entry {
7162            ts: ts(10, 0, 0),
7163            symbol: "EURUSD".into(),
7164            side: Side::Buy,
7165            order_type: OrderType::Limit,
7166            price: Some(1.0),
7167            risk_multiplier: 1.0,
7168            stoploss: None,
7169            targets: Vec::new(),
7170            group: None,
7171            trade_id: Some("cancellation-blocker".into()),
7172            entry_class: None,
7173        };
7174        let outcome = runner.run_raw_signals_controlled(
7175            &mut feed,
7176            vec![pending],
7177            None,
7178            || cancelled.get(),
7179            |progress| {
7180                if progress.processed_events >= REPLAY_PROGRESS_INTERVAL {
7181                    cancelled.set(true);
7182                }
7183            },
7184        );
7185
7186        assert_eq!(outcome.unwrap_err(), ReplayCancelled);
7187    }
7188
7189    #[test]
7190    fn controlled_replay_progress_is_monotonic_and_reaches_event_total() {
7191        let mut feed = long_tick_feed(600);
7192        let mut updates = Vec::new();
7193        BacktestRunner::with_defaults()
7194            .run_raw_signals_controlled(
7195                &mut feed,
7196                Vec::new(),
7197                None,
7198                || false,
7199                |progress| updates.push(progress),
7200            )
7201            .unwrap();
7202
7203        assert!(updates.len() >= 3);
7204        assert!(updates.windows(2).all(|pair| {
7205            pair[0].processed_events <= pair[1].processed_events
7206                && pair[0].processed_signals <= pair[1].processed_signals
7207                && pair[0].total_events <= pair[1].total_events
7208                && pair[0].total_signals <= pair[1].total_signals
7209        }));
7210        assert_eq!(updates.last().unwrap().processed_events, 600);
7211        assert_eq!(updates.last().unwrap().total_events, 600);
7212    }
7213
7214    #[test]
7215    fn legacy_replay_skips_invalid_crossed_and_reversed_quotes_without_nonfinite_pnl() {
7216        let events = vec![
7217            tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7218            tick("EURUSD", f64::NAN, 101.0, ts(10, 0, 1)),
7219            tick("EURUSD", 102.0, 101.0, ts(10, 0, 2)),
7220            tick("EURUSD", 90.0, 90.0, ts(9, 59, 59)),
7221            tick("EURUSD", 110.0, 110.0, ts(10, 0, 3)),
7222        ];
7223        let mut feed = VecFeed::new(events);
7224        let signals = vec![RawSignal::Entry {
7225            ts: ts(10, 0, 0),
7226            symbol: "EURUSD".into(),
7227            side: Side::Buy,
7228            order_type: OrderType::Market,
7229            price: Some(100.0),
7230            risk_multiplier: 1.0,
7231            stoploss: None,
7232            targets: vec![],
7233            group: None,
7234            trade_id: Some("safe-feed".into()),
7235            entry_class: None,
7236        }];
7237
7238        let result =
7239            BacktestRunner::new(fixed_lot_config()).run_raw_signals(&mut feed, signals, None);
7240        assert_eq!(result.trade_log.len(), 1);
7241        assert_eq!(result.trade_log[0].exit_price, 110.0);
7242        assert_eq!(result.trade_log[0].pnl, 10.0);
7243        assert!(result.total_pnl.is_finite());
7244        assert!(result.final_balance.is_finite());
7245    }
7246
7247    #[test]
7248    fn legacy_and_future_profile_replay_share_empty_ratio_target_resolution() {
7249        let profile = ManagementProfile {
7250            name: "equal-target".into(),
7251            target_selection: None,
7252            use_targets: vec![1],
7253            close_ratios: vec![],
7254            target_source: TargetSource::FromSignal,
7255            stoploss_mode: StoplossMode::FromSignal,
7256            rules: vec![],
7257            group_override: None,
7258            let_remainder_run: false,
7259            entry_geometry: EntryGeometryPolicy::Strict,
7260        };
7261        let signals = vec![RawSignal::Entry {
7262            ts: ts(10, 0, 0),
7263            symbol: "EURUSD".into(),
7264            side: Side::Buy,
7265            order_type: OrderType::Market,
7266            price: Some(100.0),
7267            risk_multiplier: 1.0,
7268            stoploss: None,
7269            targets: vec![101.0],
7270            group: None,
7271            trade_id: Some("profile-parity".into()),
7272            entry_class: None,
7273        }];
7274        let events = vec![
7275            tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7276            tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7277        ];
7278
7279        let mut legacy_feed = VecFeed::new(events.clone());
7280        let legacy = BacktestRunner::new(BacktestConfig {
7281            close_on_finish: false,
7282            ..fixed_lot_config()
7283        })
7284        .run_raw_signals(&mut legacy_feed, signals.clone(), Some(&profile));
7285        let mut future_feed = VecFeed::new(events);
7286        let future = BacktestRunner::new_future(
7287            BacktestConfig {
7288                close_on_finish: false,
7289                ..fixed_lot_config()
7290            },
7291            FutureQuoteConfig::default(),
7292        )
7293        .run_raw_signals_future(&mut future_feed, signals, Some(&profile));
7294
7295        assert_eq!(legacy.trade_log.len(), 1);
7296        assert_eq!(future.trade_log.len(), 1);
7297        assert_eq!(legacy.trade_log[0].close_reason, CloseReason::Target);
7298        assert_eq!(future.trade_log[0].close_reason, CloseReason::Target);
7299        assert_eq!(legacy.trade_log[0].size, future.trade_log[0].size);
7300    }
7301
7302    #[test]
7303    fn future_batch_sizes_from_shared_conversion_before_primary_and_uses_primary_eod() {
7304        let currency_plan = RunCurrencyPlan::new(
7305            "USD",
7306            ["EURUSD".to_owned()].into_iter().collect(),
7307            ["EURUSD".to_owned()].into_iter().collect(),
7308            [("EURUSD".to_owned(), "EUR".to_owned())]
7309                .into_iter()
7310                .collect(),
7311            [(
7312                "EUR".to_owned(),
7313                ConversionRoute::Direct {
7314                    pair: FxPair {
7315                        symbol: "EURUSD".to_owned(),
7316                        base_currency: "EUR".to_owned(),
7317                        quote_currency: "USD".to_owned(),
7318                    },
7319                },
7320            )]
7321            .into_iter()
7322            .collect(),
7323            Vec::new(),
7324        )
7325        .unwrap();
7326        let mut config = fixed_lot_config();
7327        config.sizing = Some(SizingPolicy::FixedRiskAmount { amount: 12.0 });
7328        let future = FutureQuoteConfig {
7329            currency_plan: Some(currency_plan),
7330            conversion_stale_after_ms: 1_000,
7331            ..FutureQuoteConfig::default()
7332        };
7333        let events = vec![
7334            FeedEvent::new(
7335                tick("EURUSD", 1.1, 1.2, ts(10, 0, 0)),
7336                EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
7337            ),
7338            FeedEvent::new(
7339                tick("EURUSD", 2.0, 2.1, ts(10, 0, 1)),
7340                EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
7341            ),
7342        ];
7343        let signals = vec![RawSignal::Entry {
7344            ts: ts(10, 0, 0),
7345            symbol: "EURUSD".into(),
7346            side: Side::Buy,
7347            order_type: OrderType::Market,
7348            price: Some(1.0),
7349            risk_multiplier: 1.0,
7350            stoploss: Some(1.19),
7351            targets: Vec::new(),
7352            group: None,
7353            trade_id: Some("shared-conversion".into()),
7354            entry_class: None,
7355        }];
7356
7357        let mut feed = VecFeed::from_feed_events(events);
7358        let result = BacktestRunner::new_future(config, future)
7359            .run_raw_signals_future(&mut feed, signals, None);
7360
7361        assert_eq!(result.recorded_fills.len(), 2);
7362        assert!((result.recorded_fills[0].fill.price - 1.2).abs() < 1.0e-12);
7363        assert!((result.recorded_fills[0].size - 10.0).abs() < 1.0e-12);
7364        assert_eq!(result.recorded_fills[1].execution_ts, Some(ts(10, 0, 0)));
7365        assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 0));
7366        assert!(
7367            result
7368                .mtm_equity_curve
7369                .iter()
7370                .all(|point| point.ts == ts(10, 0, 0))
7371        );
7372    }
7373
7374    #[test]
7375    fn conversion_only_batch_revalues_but_defers_execution_to_primary_quote() {
7376        let currency_plan = RunCurrencyPlan::new(
7377            "USD",
7378            ["EURUSD".to_owned()].into_iter().collect(),
7379            ["EURUSD".to_owned()].into_iter().collect(),
7380            [("EURUSD".to_owned(), "EUR".to_owned())]
7381                .into_iter()
7382                .collect(),
7383            [(
7384                "EUR".to_owned(),
7385                ConversionRoute::Direct {
7386                    pair: FxPair {
7387                        symbol: "EURUSD".to_owned(),
7388                        base_currency: "EUR".to_owned(),
7389                        quote_currency: "USD".to_owned(),
7390                    },
7391                },
7392            )]
7393            .into_iter()
7394            .collect(),
7395            Vec::new(),
7396        )
7397        .unwrap();
7398        let config = BacktestConfig {
7399            close_on_finish: false,
7400            ..fixed_lot_config()
7401        };
7402        let future = FutureQuoteConfig {
7403            currency_plan: Some(currency_plan),
7404            conversion_stale_after_ms: 10_000,
7405            ..FutureQuoteConfig::default()
7406        };
7407        let events = vec![
7408            FeedEvent::new(
7409                tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7410                EventMetadata::new(SeriesRoles::PRIMARY_AND_CONVERSION, 0, 0),
7411            ),
7412            FeedEvent::new(
7413                tick("EURUSD", 2.0, 2.0, ts(10, 0, 1)),
7414                EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
7415            ),
7416            FeedEvent::new(
7417                tick("EURUSD", 110.0, 110.0, ts(10, 0, 2)),
7418                EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
7419            ),
7420        ];
7421        let signals = vec![
7422            RawSignal::Entry {
7423                ts: ts(10, 0, 0),
7424                symbol: "EURUSD".into(),
7425                side: Side::Buy,
7426                order_type: OrderType::Market,
7427                price: None,
7428                risk_multiplier: 1.0,
7429                stoploss: None,
7430                targets: Vec::new(),
7431                group: None,
7432                trade_id: Some("conversion-only".into()),
7433                entry_class: None,
7434            },
7435            RawSignal::Close {
7436                ts: ts(10, 0, 1),
7437                position: PositionRef::ByTradeId {
7438                    trade_id: "conversion-only".into(),
7439                },
7440            },
7441        ];
7442
7443        let mut feed = VecFeed::from_feed_events(events);
7444        let result = BacktestRunner::new_future(config, future)
7445            .run_raw_signals_future(&mut feed, signals, None);
7446
7447        assert_eq!(result.recorded_fills.len(), 2);
7448        assert_eq!(result.recorded_fills[0].quote_ts, ts(10, 0, 0));
7449        assert_eq!(result.recorded_fills[1].quote_ts, ts(10, 0, 2));
7450        assert!(
7451            result
7452                .mtm_equity_curve
7453                .iter()
7454                .any(|point| point.ts == ts(10, 0, 1))
7455        );
7456        assert_eq!(result.total_pnl, 20.0);
7457        assert_eq!(result.close_events[0].native_pnl, Some(10.0));
7458        assert_eq!(
7459            result.close_events[0]
7460                .pnl_conversion
7461                .as_ref()
7462                .unwrap()
7463                .operation_ts,
7464            ts(10, 0, 2)
7465        );
7466    }
7467
7468    #[test]
7469    fn exact_timestamp_close_updates_balance_before_later_risk_entry() {
7470        let mut config = fixed_lot_config();
7471        config.close_on_finish = false;
7472        config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
7473        let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
7474        spec.digits = 2;
7475        spec.pip_position = 2;
7476        spec.lot_base_units = 1;
7477        spec.lot_step_units = 1;
7478        let future = FutureQuoteConfig {
7479            currency_plan: Some(identity_currency_plan("EURUSD")),
7480            ..FutureQuoteConfig::default()
7481        };
7482        let signals = vec![
7483            RawSignal::Entry {
7484                ts: ts(10, 0, 0),
7485                symbol: "EURUSD".into(),
7486                side: Side::Buy,
7487                order_type: OrderType::Market,
7488                price: None,
7489                risk_multiplier: 1.0,
7490                stoploss: Some(99.0),
7491                targets: Vec::new(),
7492                group: None,
7493                trade_id: Some("first".into()),
7494                entry_class: None,
7495            },
7496            RawSignal::Close {
7497                ts: ts(10, 0, 1),
7498                position: PositionRef::ByTradeId {
7499                    trade_id: "first".into(),
7500                },
7501            },
7502            RawSignal::Entry {
7503                ts: ts(10, 0, 1),
7504                symbol: "EURUSD".into(),
7505                side: Side::Buy,
7506                order_type: OrderType::Market,
7507                price: None,
7508                risk_multiplier: 1.0,
7509                stoploss: Some(100.0),
7510                targets: Vec::new(),
7511                group: None,
7512                trade_id: Some("second".into()),
7513                entry_class: None,
7514            },
7515        ];
7516        let mut feed = VecFeed::new(vec![
7517            tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7518            tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7519        ]);
7520
7521        let result = BacktestRunner::new_future(config, future)
7522            .run_raw_signals_future(&mut feed, signals, None);
7523
7524        assert!((result.total_pnl - 100.0).abs() < 1.0e-12);
7525        assert_eq!(result.open_position_snapshots.len(), 1);
7526        assert_eq!(
7527            result.open_position_snapshots[0].trade_id.as_deref(),
7528            Some("second")
7529        );
7530        assert!((result.open_position_snapshots[0].remaining_size - 101.0).abs() < 1.0e-12);
7531    }
7532
7533    #[test]
7534    fn pending_fill_keeps_placement_size_after_balance_changes() {
7535        let mut config = fixed_lot_config();
7536        config.close_on_finish = false;
7537        config.sizing = Some(SizingPolicy::BalanceRiskPercent { percent: 1.0 });
7538        let spec = config.symbol_specs.get_mut("EURUSD").unwrap();
7539        spec.digits = 2;
7540        spec.pip_position = 2;
7541        spec.lot_base_units = 1;
7542        spec.lot_step_units = 1;
7543        let future = FutureQuoteConfig {
7544            currency_plan: Some(identity_currency_plan("EURUSD")),
7545            market_entry_sizing_basis: MarketEntrySizingBasis::SignalEntryPrice,
7546            ..FutureQuoteConfig::default()
7547        };
7548        let signals = vec![
7549            RawSignal::Entry {
7550                ts: ts(10, 0, 0),
7551                symbol: "EURUSD".into(),
7552                side: Side::Buy,
7553                order_type: OrderType::Market,
7554                price: None,
7555                risk_multiplier: 1.0,
7556                stoploss: Some(99.0),
7557                targets: Vec::new(),
7558                group: None,
7559                trade_id: Some("market".into()),
7560                entry_class: None,
7561            },
7562            RawSignal::Entry {
7563                ts: ts(10, 0, 0),
7564                symbol: "EURUSD".into(),
7565                side: Side::Buy,
7566                order_type: OrderType::Limit,
7567                price: Some(99.0),
7568                risk_multiplier: 1.0,
7569                stoploss: Some(98.0),
7570                targets: Vec::new(),
7571                group: None,
7572                trade_id: Some("pending".into()),
7573                entry_class: None,
7574            },
7575            RawSignal::Close {
7576                ts: ts(10, 0, 1),
7577                position: PositionRef::ByTradeId {
7578                    trade_id: "market".into(),
7579                },
7580            },
7581        ];
7582        let mut feed = VecFeed::new(vec![
7583            tick("EURUSD", 100.0, 100.0, ts(10, 0, 0)),
7584            tick("EURUSD", 101.0, 101.0, ts(10, 0, 1)),
7585            tick("EURUSD", 99.0, 99.0, ts(10, 0, 2)),
7586        ]);
7587
7588        let result = BacktestRunner::new_future(config, future)
7589            .run_raw_signals_future(&mut feed, signals, None);
7590
7591        assert_eq!(result.pending_order_snapshots.len(), 0);
7592        assert_eq!(result.open_position_snapshots.len(), 1);
7593        assert_eq!(
7594            result.open_position_snapshots[0].trade_id.as_deref(),
7595            Some("pending")
7596        );
7597        assert!((result.open_position_snapshots[0].remaining_size - 100.0).abs() < 1.0e-12);
7598        let metadata = result.execution_metadata.as_ref().unwrap();
7599        assert_eq!(metadata.market_entry_sizing.len(), 1);
7600        assert_eq!(
7601            metadata.market_entry_sizing[0].trade_id.as_deref(),
7602            Some("market")
7603        );
7604        assert_eq!(metadata.entry_profile_resolutions.len(), 2);
7605        assert!(metadata.entry_profile_resolutions.iter().any(|audit| {
7606            audit.trade_id.as_deref() == Some("pending")
7607                && audit.resolution_stage == EntryResolutionStage::PendingPlacement
7608        }));
7609    }
7610
7611    #[test]
7612    fn raw_entries_require_sizing_but_management_only_replay_does_not() {
7613        let entry = RawSignal::Entry {
7614            ts: ts(10, 0, 0),
7615            symbol: "EURUSD".into(),
7616            side: Side::Buy,
7617            order_type: OrderType::Market,
7618            price: None,
7619            risk_multiplier: 1.0,
7620            stoploss: None,
7621            targets: Vec::new(),
7622            group: None,
7623            trade_id: None,
7624            entry_class: None,
7625        };
7626        let mut entry_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
7627        let rejected =
7628            BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
7629                .run_raw_signals_future(&mut entry_feed, vec![entry], None);
7630        assert!(rejected.action_dispositions.iter().any(|disposition| {
7631            disposition.action_id == "configuration"
7632                && disposition
7633                    .reason
7634                    .as_deref()
7635                    .is_some_and(|reason| reason.contains("BacktestConfig.sizing"))
7636        }));
7637
7638        let mut management_feed = VecFeed::new(vec![tick("EURUSD", 100.0, 100.0, ts(10, 0, 0))]);
7639        let management =
7640            BacktestRunner::new_future(BacktestConfig::default(), FutureQuoteConfig::default())
7641                .run_raw_signals_future(
7642                    &mut management_feed,
7643                    vec![RawSignal::CloseAll { ts: ts(10, 0, 0) }],
7644                    None,
7645                );
7646        assert!(
7647            management
7648                .action_dispositions
7649                .iter()
7650                .all(|disposition| disposition.action_id != "configuration")
7651        );
7652    }
7653
7654    #[test]
7655    fn server_filter_signals_before_market_window() {
7656        // Regression test for Issue 1 Part 4:
7657        // Verify the runner does NOT filter pre-window signals (library level).
7658        // The server filter is tested separately in handlers.
7659        // Here we verify that signals with ts before first market event
7660        // ARE still injected (library behavior). Server filtering removes them.
7661        let events = vec![
7662            tick("EURUSD", 1.0848, 1.0850, ts(10, 0, 0)),
7663            tick("EURUSD", 1.0900, 1.0902, ts(10, 0, 1)),
7664        ];
7665        let mut feed = VecFeed::new(events);
7666
7667        // Signal from January, market data from "today" (ts(10,0,0)).
7668        let raw_signals = vec![RawSignal::Entry {
7669            ts: NaiveDate::from_ymd_opt(2026, 1, 1)
7670                .unwrap()
7671                .and_hms_opt(0, 0, 0)
7672                .unwrap(),
7673            symbol: "EURUSD".into(),
7674            side: Side::Buy,
7675            order_type: OrderType::Market,
7676            price: Some(1.0850),
7677            risk_multiplier: 1.0,
7678            stoploss: None,
7679            targets: vec![1.0900],
7680            group: None,
7681            trade_id: None,
7682            entry_class: None,
7683        }];
7684
7685        let runner = BacktestRunner::new(fixed_lot_config());
7686        let result = runner.run_raw_signals(&mut feed, raw_signals, None);
7687
7688        // Library still injects it; server filtering is the authoritative gate.
7689        assert_eq!(result.total_trades, 1);
7690    }
7691
7692    #[test]
7693    fn future_replay_audits_profile_resolution_rejection() {
7694        let profile = ManagementProfile {
7695            name: "requires_stop".into(),
7696            target_selection: Some(crate::profile::TargetSelection::None),
7697            use_targets: vec![],
7698            close_ratios: vec![],
7699            target_source: TargetSource::FromSignal,
7700            stoploss_mode: StoplossMode::FromSignalDistance { multiplier: 1.5 },
7701            rules: vec![],
7702            group_override: None,
7703            let_remainder_run: true,
7704            entry_geometry: EntryGeometryPolicy::Strict,
7705        };
7706        let profiles = PreparedEntryProfiles::try_new(
7707            Some(profile),
7708            Vec::<(String, ManagementProfile)>::new(),
7709        )
7710        .unwrap();
7711        let signals = vec![RawSignal::Entry {
7712            ts: ts(10, 0, 0),
7713            symbol: "EURUSD".into(),
7714            side: Side::Buy,
7715            order_type: OrderType::Market,
7716            price: Some(1.1000),
7717            risk_multiplier: 1.0,
7718            stoploss: None,
7719            targets: vec![],
7720            group: None,
7721            trade_id: Some("missing-stop".into()),
7722            entry_class: None,
7723        }];
7724        let mut feed = VecFeed::new(vec![tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1))]);
7725        let result = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default())
7726            .with_entry_profiles(profiles)
7727            .run_raw_signals_future(&mut feed, signals, None);
7728        let audit = &result
7729            .execution_metadata
7730            .as_ref()
7731            .unwrap()
7732            .entry_profile_resolutions[0];
7733        assert_eq!(audit.outcome, ActionDispositionStatus::Rejected);
7734        assert_eq!(audit.rejection_stage.as_deref(), Some("profile_resolution"));
7735        assert!(audit.reason.as_deref().unwrap().contains("signal stoploss"));
7736    }
7737
7738    #[test]
7739    fn future_replay_routes_entry_profiles_and_audits_resolved_levels() {
7740        let default_profile = ManagementProfile {
7741            name: "default".into(),
7742            target_selection: Some(crate::profile::TargetSelection::None),
7743            use_targets: vec![],
7744            close_ratios: vec![],
7745            target_source: TargetSource::FromSignal,
7746            stoploss_mode: StoplossMode::FromSignal,
7747            rules: vec![],
7748            group_override: None,
7749            let_remainder_run: true,
7750            entry_geometry: EntryGeometryPolicy::Strict,
7751        };
7752        let expanded_profile = ManagementProfile {
7753            name: "expanded".into(),
7754            target_selection: None,
7755            use_targets: vec![],
7756            close_ratios: vec![1.0],
7757            target_source: TargetSource::StopDistanceMultiples {
7758                multiples: vec![1.0],
7759            },
7760            stoploss_mode: StoplossMode::FromSignalDistance { multiplier: 1.5 },
7761            rules: vec![],
7762            group_override: None,
7763            let_remainder_run: false,
7764            entry_geometry: EntryGeometryPolicy::Strict,
7765        };
7766        let profiles = PreparedEntryProfiles::try_new(
7767            Some(default_profile),
7768            [("expanded".to_owned(), expanded_profile)],
7769        )
7770        .unwrap();
7771        let signals = vec![
7772            RawSignal::Entry {
7773                ts: ts(10, 0, 0),
7774                symbol: "EURUSD".into(),
7775                side: Side::Buy,
7776                order_type: OrderType::Market,
7777                price: Some(1.1000),
7778                risk_multiplier: 1.0,
7779                stoploss: Some(1.0990),
7780                targets: vec![],
7781                group: Some("same-group".into()),
7782                trade_id: Some("default-entry".into()),
7783                entry_class: None,
7784            },
7785            RawSignal::Entry {
7786                ts: ts(10, 0, 1),
7787                symbol: "EURUSD".into(),
7788                side: Side::Buy,
7789                order_type: OrderType::Market,
7790                price: Some(1.1000),
7791                risk_multiplier: 1.0,
7792                stoploss: Some(1.0990),
7793                targets: vec![],
7794                group: Some("same-group".into()),
7795                trade_id: Some("expanded-entry".into()),
7796                entry_class: Some("expanded".into()),
7797            },
7798        ];
7799        let mut feed = VecFeed::new(vec![
7800            tick("EURUSD", 1.1000, 1.1000, ts(10, 0, 1)),
7801            tick("EURUSD", 1.1002, 1.1002, ts(10, 0, 2)),
7802            tick("EURUSD", 1.1003, 1.1003, ts(10, 0, 3)),
7803        ]);
7804        let result = BacktestRunner::new_future(fixed_lot_config(), FutureQuoteConfig::default())
7805            .with_entry_profiles(profiles)
7806            .run_raw_signals_future(&mut feed, signals, None);
7807
7808        let audits = &result
7809            .execution_metadata
7810            .as_ref()
7811            .unwrap()
7812            .entry_profile_resolutions;
7813        assert_eq!(audits.len(), 2);
7814        assert_eq!(
7815            audits[0].selection_source,
7816            EntryProfileSelectionSource::RunDefault
7817        );
7818        assert_eq!(audits[0].selected_profile_name.as_deref(), Some("default"));
7819        assert_eq!(
7820            audits[1].selection_source,
7821            EntryProfileSelectionSource::Mapped
7822        );
7823        assert_eq!(audits[1].entry_class.as_deref(), Some("expanded"));
7824        assert_eq!(audits[1].selected_profile_name.as_deref(), Some("expanded"));
7825        let levels = audits[1].level_resolution.as_ref().unwrap();
7826        let reference = audits[1].level_reference_price.unwrap();
7827        let stop = levels.resolved_stoploss.unwrap();
7828        let target = levels.resolved_targets[0];
7829        let risk_distance = reference - stop;
7830        assert!(risk_distance > 0.0);
7831        assert!((target - reference - risk_distance).abs() < 1.0e-9);
7832    }
7833}