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