Skip to main content

qs_backtest/
future_executor.rs

1//! Fill-authoritative accounting for the FutureQuoteV1 replay path.
2
3use std::collections::{BTreeMap, BTreeSet, HashMap};
4
5use chrono::NaiveDateTime;
6
7use qs_core::TradeEngine;
8use qs_core::types::{
9    CloseReason, Effect, EffectiveStop, FillPurpose, FutureEffect, FutureFill, PriceQuote, Side,
10    StopOrigin, position_size_tolerance,
11};
12use thiserror::Error;
13
14use crate::artifacts::{
15    CloseEvent, CompletedPosition, OpenPositionSnapshot, PendingOrderLifecycleEvent,
16    PendingOrderLifecycleState, RecordedFill, RiskBasisStatus, RiskTranche, deterministic_event_id,
17};
18use crate::currency::{
19    ConversionError, ConversionQuoteBook, ConversionResult, ConversionRoute, RunCurrencyPlan,
20};
21use crate::portfolio::PortfolioRecorder;
22use crate::report::TradeResult;
23
24#[derive(Debug, Clone)]
25struct PendingOrigin {
26    position_id: String,
27    placement_action_id: Option<String>,
28    signal_ts: Option<NaiveDateTime>,
29    effective_ts: NaiveDateTime,
30    placed_ts: NaiveDateTime,
31    symbol: String,
32    side: Side,
33    order_type: qs_core::types::OrderType,
34    requested_size: f64,
35    requested_price: Option<f64>,
36    placement_sequence: u64,
37    initial_stop: Option<f64>,
38}
39
40#[derive(Debug, Error)]
41pub(crate) enum FutureExecutorError {
42    #[error("position not found while processing FutureQuote effect: {0}")]
43    PositionNotFound(String),
44    #[error("account not found while processing FutureQuote effect: {0}")]
45    AccountNotFound(String),
46    #[error("fill-bearing effect was emitted without a fill: {0}")]
47    MissingFill(String),
48    #[error("non-fill effect unexpectedly carried a fill: {0}")]
49    UnexpectedFill(String),
50    #[error("invalid carried fill for {position_id}: {reason}")]
51    InvalidFill { position_id: String, reason: String },
52    #[error("portfolio rejected realized P&L {pnl} for {position_id}")]
53    PortfolioRejectedRealizedPnl { position_id: String, pnl: f64 },
54    #[error("currency plan has no P&L currency or route for primary symbol {0}")]
55    MissingCurrencyRoute(String),
56
57    #[error("account conversion failed for {symbol} at {operation_ts}: {source}")]
58    Conversion {
59        symbol: String,
60        operation_ts: NaiveDateTime,
61        #[source]
62        source: ConversionError,
63    },
64    #[error("account conversion for {symbol} produced invalid {kind} amount {amount}")]
65    InvalidConvertedAmount {
66        symbol: String,
67        kind: &'static str,
68        amount: f64,
69    },
70}
71
72#[derive(Debug, Clone)]
73struct PositionAccount {
74    position_id: String,
75    symbol: String,
76    side: Side,
77    group: Option<String>,
78    trade_id: Option<String>,
79    open_ts: NaiveDateTime,
80    entry_size: f64,
81    remaining_size: f64,
82    /// Value of every historical entry, used only for campaign-level audit.
83    entry_value: f64,
84    /// Average-cost basis assigned only to inventory that is still open.
85    open_entry_value: f64,
86    initial_stop: Option<f64>,
87    effective_stop: Option<EffectiveStop>,
88    risk_tranches: Vec<RiskTranche>,
89    close_events: Vec<CloseEvent>,
90    realized_pnl: f64,
91    native_realized_pnl: f64,
92    native_currency: Option<String>,
93    account_currency: Option<String>,
94}
95
96#[derive(Debug, Clone)]
97struct AccountedAmount {
98    amount: f64,
99    native_currency: Option<String>,
100    conversion: Option<ConversionResult>,
101}
102
103impl PositionAccount {
104    fn average_entry(&self) -> f64 {
105        if self.remaining_size <= position_size_tolerance(self.entry_size) {
106            0.0
107        } else {
108            self.open_entry_value / self.remaining_size
109        }
110    }
111
112    fn historical_average_entry(&self) -> f64 {
113        if self.entry_size <= position_size_tolerance(self.entry_size) {
114            0.0
115        } else {
116            self.entry_value / self.entry_size
117        }
118    }
119
120    fn snapshot(&self) -> OpenPositionSnapshot {
121        let mut snapshot = OpenPositionSnapshot::new(
122            self.position_id.clone(),
123            self.symbol.clone(),
124            self.side,
125            self.average_entry(),
126            self.remaining_size,
127        );
128        snapshot.group = self.group.clone();
129        snapshot.trade_id = self.trade_id.clone();
130        snapshot.open_ts = Some(self.open_ts);
131        snapshot.initial_stop = self.initial_stop;
132        snapshot.effective_stop = self.effective_stop;
133        snapshot.realized_pnl = self.realized_pnl;
134        snapshot.native_realized_pnl = Some(self.native_realized_pnl);
135        snapshot.native_currency = self.native_currency.clone();
136        snapshot.account_currency = self.account_currency.clone();
137        snapshot
138    }
139}
140
141/// Accounting adapter for quote-authoritative FutureQuoteV1 execution.
142#[derive(Debug, Clone)]
143pub struct FutureExecutor {
144    initial_balance: f64,
145    balance: f64,
146    contract_sizes: HashMap<String, f64>,
147    currency_plan: Option<RunCurrencyPlan>,
148    accounts: BTreeMap<String, PositionAccount>,
149    pending_origins: BTreeMap<String, PendingOrigin>,
150    terminal_pending_orders: BTreeSet<String>,
151
152    pub fills: Vec<RecordedFill>,
153    pub pending_order_lifecycle: Vec<PendingOrderLifecycleEvent>,
154    pub close_events: Vec<CloseEvent>,
155    pub completed_positions: Vec<CompletedPosition>,
156    pub trade_log: Vec<TradeResult>,
157    fill_sequence: u64,
158    close_sequence: u64,
159    pending_lifecycle_sequence: u64,
160    pnl_epsilon: f64,
161}
162
163#[derive(Debug)]
164struct FutureExecutorCheckpoint {
165    accounts: BTreeMap<String, Option<PositionAccount>>,
166    pending_origins: BTreeMap<String, Option<PendingOrigin>>,
167    terminal_pending_orders: BTreeMap<String, bool>,
168    balance: f64,
169    fill_sequence: u64,
170    close_sequence: u64,
171    pending_lifecycle_sequence: u64,
172    fills_len: usize,
173    pending_order_lifecycle_len: usize,
174    close_events_len: usize,
175    completed_positions_len: usize,
176    trade_log_len: usize,
177}
178
179#[derive(Debug)]
180struct PendingCampaignCompletion {
181    position_id: String,
182    final_net_pnl: f64,
183    completed_position_index: usize,
184}
185
186#[derive(Debug)]
187struct PortfolioBatch {
188    realized_pnl: f64,
189    realized_pnl_changed: bool,
190    campaign_completions: Vec<PendingCampaignCompletion>,
191}
192
193impl PortfolioBatch {
194    fn new(portfolio: &PortfolioRecorder) -> Self {
195        Self {
196            realized_pnl: portfolio.realized_pnl(),
197            realized_pnl_changed: false,
198            campaign_completions: Vec::new(),
199        }
200    }
201
202    fn add_realized_pnl(&mut self, pnl: f64) -> bool {
203        let next = self.realized_pnl + pnl;
204        if !pnl.is_finite() || !next.is_finite() {
205            return false;
206        }
207        self.realized_pnl = next;
208        self.realized_pnl_changed = true;
209        true
210    }
211
212    fn finish_campaign(
213        &mut self,
214        position_id: String,
215        final_net_pnl: f64,
216        completed_position_index: usize,
217    ) {
218        self.campaign_completions.push(PendingCampaignCompletion {
219            position_id,
220            final_net_pnl,
221            completed_position_index,
222        });
223    }
224
225    fn commit(self, portfolio: &mut PortfolioRecorder, completed: &mut [CompletedPosition]) {
226        if self.realized_pnl_changed {
227            let updated = portfolio.set_realized_pnl(self.realized_pnl);
228            debug_assert!(updated, "staged realized P&L was prevalidated");
229        }
230        for completion in self.campaign_completions {
231            let excursion =
232                portfolio.finish_campaign(&completion.position_id, completion.final_net_pnl);
233            if let Some(position) = completed.get_mut(completion.completed_position_index) {
234                position.mae = excursion.map(|value| value.mae);
235                position.mfe = excursion.map(|value| value.mfe);
236            }
237        }
238    }
239}
240
241impl FutureExecutorCheckpoint {
242    fn capture(executor: &FutureExecutor, effects: &[FutureEffect]) -> Self {
243        let affected_ids: BTreeSet<_> = effects
244            .iter()
245            .map(|effect| effect_position_id(effect.effect()).to_owned())
246            .collect();
247        let accounts = affected_ids
248            .iter()
249            .map(|id| (id.clone(), executor.accounts.get(id).cloned()))
250            .collect();
251        let pending_origins = affected_ids
252            .iter()
253            .map(|id| (id.clone(), executor.pending_origins.get(id).cloned()))
254            .collect();
255        let terminal_pending_orders = affected_ids
256            .into_iter()
257            .map(|id| {
258                let present = executor.terminal_pending_orders.contains(&id);
259                (id, present)
260            })
261            .collect();
262        Self {
263            accounts,
264            pending_origins,
265            terminal_pending_orders,
266            balance: executor.balance,
267            fill_sequence: executor.fill_sequence,
268            close_sequence: executor.close_sequence,
269            pending_lifecycle_sequence: executor.pending_lifecycle_sequence,
270            fills_len: executor.fills.len(),
271            pending_order_lifecycle_len: executor.pending_order_lifecycle.len(),
272            close_events_len: executor.close_events.len(),
273            completed_positions_len: executor.completed_positions.len(),
274            trade_log_len: executor.trade_log.len(),
275        }
276    }
277
278    fn restore(self, executor: &mut FutureExecutor) {
279        restore_entries(&mut executor.accounts, self.accounts);
280        restore_entries(&mut executor.pending_origins, self.pending_origins);
281        for (id, present) in self.terminal_pending_orders {
282            if present {
283                executor.terminal_pending_orders.insert(id);
284            } else {
285                executor.terminal_pending_orders.remove(&id);
286            }
287        }
288        executor.balance = self.balance;
289        executor.fill_sequence = self.fill_sequence;
290        executor.close_sequence = self.close_sequence;
291        executor.pending_lifecycle_sequence = self.pending_lifecycle_sequence;
292        executor.fills.truncate(self.fills_len);
293        executor
294            .pending_order_lifecycle
295            .truncate(self.pending_order_lifecycle_len);
296        executor.close_events.truncate(self.close_events_len);
297        executor
298            .completed_positions
299            .truncate(self.completed_positions_len);
300        executor.trade_log.truncate(self.trade_log_len);
301    }
302}
303
304fn restore_entries<T>(entries: &mut BTreeMap<String, T>, checkpoint: BTreeMap<String, Option<T>>) {
305    for (id, value) in checkpoint {
306        match value {
307            Some(value) => {
308                entries.insert(id, value);
309            }
310            None => {
311                entries.remove(&id);
312            }
313        }
314    }
315}
316
317fn effect_position_id(effect: &Effect) -> &str {
318    match effect {
319        Effect::OrderPlaced { id }
320        | Effect::OrderCancelled { id }
321        | Effect::PositionOpened { id }
322        | Effect::PositionClosed { id, .. }
323        | Effect::PartialClose { id, .. }
324        | Effect::StoplossModified { id, .. }
325        | Effect::StoplossRemoved { id, .. }
326        | Effect::ScaledIn { id, .. }
327        | Effect::RuleTriggered { id, .. } => id,
328    }
329}
330
331impl FutureExecutor {
332    pub fn new(
333        initial_balance: f64,
334        contract_sizes: HashMap<String, f64>,
335        pnl_epsilon: f64,
336    ) -> Self {
337        Self {
338            initial_balance,
339            balance: initial_balance,
340            contract_sizes,
341            currency_plan: None,
342            accounts: BTreeMap::new(),
343            pending_origins: BTreeMap::new(),
344            terminal_pending_orders: BTreeSet::new(),
345
346            fills: Vec::new(),
347            pending_order_lifecycle: Vec::new(),
348            close_events: Vec::new(),
349            completed_positions: Vec::new(),
350            trade_log: Vec::new(),
351            fill_sequence: 0,
352            close_sequence: 0,
353            pending_lifecycle_sequence: 0,
354            pnl_epsilon: if pnl_epsilon.is_finite() {
355                pnl_epsilon.abs()
356            } else {
357                1.0e-9
358            },
359        }
360    }
361
362    pub fn with_currency_plan(mut self, currency_plan: Option<RunCurrencyPlan>) -> Self {
363        self.currency_plan = currency_plan;
364        self
365    }
366
367    pub fn balance(&self) -> f64 {
368        self.balance
369    }
370
371    pub fn realized_pnl(&self) -> f64 {
372        self.balance - self.initial_balance
373    }
374
375    pub(crate) fn requires_processing(effects: &[FutureEffect]) -> bool {
376        effects.iter().any(|effect| {
377            !matches!(
378                effect,
379                FutureEffect::Plain {
380                    effect: Effect::RuleTriggered { .. },
381                    ..
382                }
383            )
384        })
385    }
386
387    pub fn has_close(&self, position_id: &str) -> bool {
388        self.accounts
389            .get(position_id)
390            .is_some_and(|account| !account.close_events.is_empty())
391            || self
392                .completed_positions
393                .iter()
394                .any(|position| position.position_id == position_id)
395    }
396
397    pub fn pending_metadata(
398        &self,
399        position_id: &str,
400    ) -> Option<(String, NaiveDateTime, NaiveDateTime)> {
401        self.pending_origins.get(position_id).and_then(|origin| {
402            Some((
403                origin.placement_action_id.clone()?,
404                origin.signal_ts?,
405                origin.effective_ts,
406            ))
407        })
408    }
409
410    pub fn open_snapshots(&self) -> Vec<OpenPositionSnapshot> {
411        self.accounts
412            .values()
413            .filter(|account| account.remaining_size > position_size_tolerance(account.entry_size))
414            .map(PositionAccount::snapshot)
415            .collect()
416    }
417
418    /// Mark every still-pending order as unfilled at the deterministic end of
419    /// the accepted quote stream. Pending snapshots remain available separately.
420    pub fn finalize_pending_orders_at_end(&mut self, terminal_ts: NaiveDateTime) {
421        let mut pending: Vec<_> = self.pending_origins.values().cloned().collect();
422        pending.sort_by_key(|origin| origin.placement_sequence);
423        for origin in pending {
424            self.record_pending_terminal(
425                &origin,
426                PendingOrderLifecycleState::UnfilledAtEnd,
427                None,
428                0.0,
429                None,
430                terminal_ts,
431            );
432        }
433    }
434
435    #[allow(dead_code, clippy::too_many_arguments)]
436    pub(crate) fn process_future_effects(
437        &mut self,
438        effects: &[FutureEffect],
439        engine: &TradeEngine,
440        quote: &PriceQuote,
441        action_id: Option<&str>,
442        signal_ts: Option<NaiveDateTime>,
443        effective_ts: NaiveDateTime,
444        portfolio: &mut PortfolioRecorder,
445    ) -> Result<Vec<String>, FutureExecutorError> {
446        self.process_future_effects_with_currency(
447            effects,
448            engine,
449            quote,
450            action_id,
451            signal_ts,
452            effective_ts,
453            portfolio,
454            None,
455        )
456    }
457
458    #[allow(clippy::too_many_arguments)]
459    pub(crate) fn process_future_effects_with_currency(
460        &mut self,
461        effects: &[FutureEffect],
462        engine: &TradeEngine,
463        quote: &PriceQuote,
464        action_id: Option<&str>,
465        signal_ts: Option<NaiveDateTime>,
466        effective_ts: NaiveDateTime,
467        portfolio: &mut PortfolioRecorder,
468        conversion_quotes: Option<&ConversionQuoteBook>,
469    ) -> Result<Vec<String>, FutureExecutorError> {
470        if !Self::requires_processing(effects) {
471            return Ok(Vec::new());
472        }
473
474        let checkpoint = FutureExecutorCheckpoint::capture(self, effects);
475        let mut portfolio_batch = PortfolioBatch::new(portfolio);
476        let result = (|| -> Result<Vec<String>, FutureExecutorError> {
477            let mut affected = Vec::new();
478            for future_effect in effects {
479                match future_effect {
480                    FutureEffect::Plain {
481                        effect,
482                        stop_origin,
483                        ..
484                    } => match effect {
485                        Effect::PositionOpened { .. }
486                        | Effect::PositionClosed { .. }
487                        | Effect::PartialClose { .. }
488                        | Effect::ScaledIn { .. } => {
489                            return Err(FutureExecutorError::MissingFill(format!("{effect:?}")));
490                        }
491                        Effect::OrderPlaced { id } => {
492                            let position = engine
493                                .get_position(id)
494                                .ok_or_else(|| FutureExecutorError::PositionNotFound(id.clone()))?;
495                            let origin = PendingOrigin {
496                                position_id: id.clone(),
497                                placement_action_id: action_id.map(str::to_owned),
498                                signal_ts,
499                                effective_ts,
500                                placed_ts: quote.ts,
501                                symbol: position.data.symbol.clone(),
502                                side: position.data.side,
503                                order_type: position.data.order_type,
504                                requested_size: position.data.size,
505                                requested_price: position.data.pending_price,
506                                placement_sequence: self.pending_lifecycle_sequence,
507                                initial_stop: position.current_stoploss(),
508                            };
509                            self.record_pending_placed(&origin);
510                            self.pending_origins.insert(id.clone(), origin);
511                        }
512                        Effect::OrderCancelled { id } => {
513                            if let Some(origin) = self.pending_origins.remove(id) {
514                                self.record_pending_terminal(
515                                    &origin,
516                                    PendingOrderLifecycleState::Cancelled,
517                                    action_id.map(str::to_owned),
518                                    0.0,
519                                    None,
520                                    quote.ts,
521                                );
522                            }
523                        }
524                        Effect::StoplossModified { id, new_price, .. } => {
525                            if !new_price.is_finite() || *new_price <= 0.0 {
526                                return Err(FutureExecutorError::InvalidFill {
527                                    position_id: id.clone(),
528                                    reason: format!(
529                                        "stoploss must be finite and positive, got {new_price}"
530                                    ),
531                                });
532                            }
533                            if let Some(account) = self.accounts.get_mut(id) {
534                                account.effective_stop = Some(EffectiveStop::new(
535                                    *new_price,
536                                    stop_origin.unwrap_or(StopOrigin::Modified),
537                                ));
538                            } else if self.pending_origins.contains_key(id) {
539                                // The immutable initial stop stays on the pending origin;
540                                // the engine supplies the effective stop when it fills.
541                            } else {
542                                return Err(FutureExecutorError::AccountNotFound(id.clone()));
543                            }
544                            affected.push(id.clone());
545                        }
546                        Effect::StoplossRemoved { id, .. } => {
547                            if let Some(account) = self.accounts.get_mut(id) {
548                                account.effective_stop = None;
549                            } else if self.pending_origins.contains_key(id) {
550                                // Removing a pending stop changes engine state only; the
551                                // placement-time initial stop remains an audit field.
552                            } else {
553                                return Err(FutureExecutorError::AccountNotFound(id.clone()));
554                            }
555                            affected.push(id.clone());
556                        }
557                        Effect::RuleTriggered { .. } => {}
558                    },
559                    FutureEffect::Filled { effect, fill, .. } => {
560                        self.validate_carried_fill(effect, fill, quote)?;
561                        match effect {
562                            Effect::PositionOpened { id } => {
563                                self.record_open(
564                                    id,
565                                    fill,
566                                    engine,
567                                    quote,
568                                    action_id,
569                                    signal_ts,
570                                    effective_ts,
571                                    conversion_quotes,
572                                )?;
573                                affected.push(id.clone());
574                            }
575                            Effect::ScaledIn { id, .. } => {
576                                self.record_scale_in(
577                                    id,
578                                    fill,
579                                    engine,
580                                    quote,
581                                    action_id,
582                                    signal_ts,
583                                    effective_ts,
584                                    conversion_quotes,
585                                )?;
586                                affected.push(id.clone());
587                            }
588                            Effect::PositionClosed { id, reason } => {
589                                self.record_close(
590                                    id,
591                                    *reason,
592                                    fill,
593                                    quote,
594                                    action_id,
595                                    signal_ts,
596                                    effective_ts,
597                                    &mut portfolio_batch,
598                                    true,
599                                    conversion_quotes,
600                                )?;
601                                if self.accounts.contains_key(id) {
602                                    return Err(FutureExecutorError::InvalidFill {
603                                        position_id: id.clone(),
604                                        reason: "full-close effect left an open account".into(),
605                                    });
606                                }
607                                affected.push(id.clone());
608                            }
609                            Effect::PartialClose { id, reason, .. } => {
610                                self.record_close(
611                                    id,
612                                    *reason,
613                                    fill,
614                                    quote,
615                                    action_id,
616                                    signal_ts,
617                                    effective_ts,
618                                    &mut portfolio_batch,
619                                    false,
620                                    conversion_quotes,
621                                )?;
622                                affected.push(id.clone());
623                            }
624                            _ => {
625                                return Err(FutureExecutorError::UnexpectedFill(format!(
626                                    "{effect:?}"
627                                )));
628                            }
629                        }
630                    }
631                }
632            }
633            affected.sort();
634            affected.dedup();
635            Ok(affected)
636        })();
637
638        match result {
639            Ok(affected) => {
640                portfolio_batch.commit(portfolio, &mut self.completed_positions);
641                Ok(affected)
642            }
643            Err(error) => {
644                checkpoint.restore(self);
645                Err(error)
646            }
647        }
648    }
649
650    fn record_pending_placed(&mut self, origin: &PendingOrigin) {
651        let event = self.pending_lifecycle_event(
652            origin,
653            PendingOrderLifecycleState::Placed,
654            None,
655            None,
656            None,
657            None,
658        );
659        self.pending_order_lifecycle.push(event);
660    }
661
662    fn record_pending_terminal(
663        &mut self,
664        origin: &PendingOrigin,
665        state: PendingOrderLifecycleState,
666        terminal_action_id: Option<String>,
667        filled_size: f64,
668        fill_price: Option<f64>,
669        terminal_ts: NaiveDateTime,
670    ) {
671        debug_assert!(state.is_terminal());
672        if !self
673            .terminal_pending_orders
674            .insert(origin.position_id.clone())
675        {
676            return;
677        }
678        let event = self.pending_lifecycle_event(
679            origin,
680            state,
681            terminal_action_id,
682            Some(filled_size),
683            fill_price,
684            Some(terminal_ts),
685        );
686        self.pending_order_lifecycle.push(event);
687    }
688
689    fn pending_lifecycle_event(
690        &mut self,
691        origin: &PendingOrigin,
692        state: PendingOrderLifecycleState,
693        terminal_action_id: Option<String>,
694        filled_size: Option<f64>,
695        fill_price: Option<f64>,
696        terminal_ts: Option<NaiveDateTime>,
697    ) -> PendingOrderLifecycleEvent {
698        let sequence = self.pending_lifecycle_sequence;
699        self.pending_lifecycle_sequence += 1;
700        let kind = match state {
701            PendingOrderLifecycleState::Placed => "pending_placed",
702            PendingOrderLifecycleState::Filled => "pending_filled",
703            PendingOrderLifecycleState::Cancelled => "pending_cancelled",
704            PendingOrderLifecycleState::UnfilledAtEnd => "pending_unfilled_at_end",
705        };
706        let wait_latency_ms = terminal_ts
707            .map(|terminal_ts| (terminal_ts - origin.placed_ts).num_milliseconds().max(0));
708        let fill_ratio = filled_size.and_then(|filled_size| {
709            (origin.requested_size.is_finite() && origin.requested_size > 0.0)
710                .then_some(filled_size / origin.requested_size)
711        });
712
713        PendingOrderLifecycleEvent {
714            id: deterministic_event_id(&origin.position_id, kind, sequence),
715            sequence,
716            position_id: origin.position_id.clone(),
717            placement_action_id: origin.placement_action_id.clone(),
718            terminal_action_id,
719            state,
720            symbol: origin.symbol.clone(),
721            side: origin.side,
722            order_type: origin.order_type,
723            requested_size: origin.requested_size,
724            filled_size,
725            requested_price: origin.requested_price,
726            fill_price,
727            signal_ts: origin.signal_ts,
728            placed_ts: Some(origin.placed_ts),
729            effective_ts: Some(origin.effective_ts),
730            terminal_ts,
731            wait_latency_ms,
732            fill_ratio,
733        }
734    }
735
736    fn validate_carried_fill(
737        &self,
738        effect: &Effect,
739        fill: &FutureFill,
740        quote: &PriceQuote,
741    ) -> Result<(), FutureExecutorError> {
742        let position_id = match effect {
743            Effect::PositionOpened { id }
744            | Effect::PositionClosed { id, .. }
745            | Effect::PartialClose { id, .. }
746            | Effect::ScaledIn { id, .. } => id.clone(),
747            _ => "<non-fill-effect>".into(),
748        };
749        if fill.source_quote_ts() != quote.ts {
750            return Err(FutureExecutorError::InvalidFill {
751                position_id,
752                reason: format!(
753                    "fill source quote timestamp {} does not match quote {}",
754                    fill.source_quote_ts(),
755                    quote.ts
756                ),
757            });
758        }
759        if fill.ts < quote.ts {
760            return Err(FutureExecutorError::InvalidFill {
761                position_id,
762                reason: format!(
763                    "fill execution timestamp {} precedes source quote {}",
764                    fill.ts, quote.ts
765                ),
766            });
767        }
768        if !fill.size.is_finite() || fill.size <= position_size_tolerance(fill.size) {
769            return Err(FutureExecutorError::InvalidFill {
770                position_id,
771                reason: format!(
772                    "size must be finite and greater than the accounting tolerance, got {}",
773                    fill.size
774                ),
775            });
776        }
777        if !fill.execution.price.is_finite() || fill.execution.price <= 0.0 {
778            return Err(FutureExecutorError::InvalidFill {
779                position_id,
780                reason: format!(
781                    "price must be finite and positive, got {}",
782                    fill.execution.price
783                ),
784            });
785        }
786        let purpose_matches = match effect {
787            Effect::PositionOpened { .. } => matches!(
788                fill.execution.purpose,
789                FillPurpose::MarketEntry | FillPurpose::LimitEntry | FillPurpose::StopEntry
790            ),
791            Effect::ScaledIn { .. } => fill.execution.purpose == FillPurpose::MarketEntry,
792            Effect::PositionClosed { reason, .. } | Effect::PartialClose { reason, .. } => {
793                fill.execution.purpose
794                    == match reason {
795                        CloseReason::Target => FillPurpose::TakeProfit,
796                        CloseReason::Stoploss
797                        | CloseReason::TrailingStop
798                        | CloseReason::BreakevenStop => FillPurpose::StopLoss,
799                        _ => FillPurpose::MarketExit,
800                    }
801            }
802            _ => false,
803        };
804        if !purpose_matches {
805            return Err(FutureExecutorError::InvalidFill {
806                position_id,
807                reason: format!(
808                    "execution purpose {:?} does not match effect",
809                    fill.execution.purpose
810                ),
811            });
812        }
813        Ok(())
814    }
815
816    #[allow(clippy::too_many_arguments)]
817    fn record_open(
818        &mut self,
819        id: &str,
820        fill: &FutureFill,
821        engine: &TradeEngine,
822        quote: &PriceQuote,
823        action_id: Option<&str>,
824        signal_ts: Option<NaiveDateTime>,
825        effective_ts: NaiveDateTime,
826        conversion_quotes: Option<&ConversionQuoteBook>,
827    ) -> Result<(), FutureExecutorError> {
828        if self.accounts.contains_key(id) {
829            return Err(FutureExecutorError::InvalidFill {
830                position_id: id.to_owned(),
831                reason: "position already has an open accounting record".into(),
832            });
833        }
834        let position = engine
835            .get_position(id)
836            .ok_or_else(|| FutureExecutorError::PositionNotFound(id.to_owned()))?;
837        if position.data.status != qs_core::types::PositionStatus::Open {
838            return Err(FutureExecutorError::InvalidFill {
839                position_id: id.to_owned(),
840                reason: format!("position is not open: {}", position.data.status),
841            });
842        }
843        if position.data.side != fill.execution.side {
844            return Err(FutureExecutorError::InvalidFill {
845                position_id: id.to_owned(),
846                reason: "execution side does not match position".into(),
847            });
848        }
849        let execution = fill.execution;
850        let size = fill.size;
851        let current_stop = position.current_effective_stop();
852        let group = position.data.group.clone();
853        let trade_id = position.data.trade_id.clone();
854        let side = position.data.side;
855        let symbol = position.data.symbol.clone();
856        let open_ts = fill.ts;
857
858        let origin = self.pending_origins.get(id).cloned();
859        let initial_stop = origin.as_ref().map_or_else(
860            || current_stop.map(|stop| stop.price),
861            |value| value.initial_stop,
862        );
863        let recorded_action_id = action_id.map(str::to_owned).or_else(|| {
864            origin
865                .as_ref()
866                .and_then(|value| value.placement_action_id.clone())
867        });
868        let recorded_signal_ts =
869            signal_ts.or_else(|| origin.as_ref().and_then(|value| value.signal_ts));
870        let recorded_effective_ts = origin
871            .as_ref()
872            .map(|value| value.effective_ts)
873            .unwrap_or(effective_ts);
874        let recorded = RecordedFill::from_quote_at(
875            id.to_owned(),
876            recorded_action_id,
877            self.fill_sequence,
878            recorded_signal_ts,
879            recorded_effective_ts,
880            fill.ts,
881            size,
882            quote,
883            execution,
884        );
885        let contract_size = self.contract_size(&symbol);
886        let risk = self.account_risk_tranche(
887            &symbol,
888            fill.ts,
889            RiskTranche::calculate(
890                Some(recorded.id.clone()),
891                side,
892                size,
893                execution.price,
894                current_stop.map(|stop| stop.price),
895                contract_size,
896                self.pnl_epsilon,
897            ),
898            conversion_quotes,
899        )?;
900        self.fill_sequence += 1;
901        self.pending_origins.remove(id);
902        self.fills.push(recorded);
903        self.accounts.insert(
904            id.to_owned(),
905            PositionAccount {
906                position_id: id.to_owned(),
907                symbol,
908                side,
909                group,
910                trade_id,
911                open_ts,
912                entry_size: size,
913                remaining_size: size,
914                entry_value: execution.price * size,
915                open_entry_value: execution.price * size,
916                initial_stop,
917                effective_stop: current_stop,
918                native_currency: risk.native_currency.clone(),
919                account_currency: self
920                    .currency_plan
921                    .as_ref()
922                    .map(|plan| plan.account_currency().to_owned()),
923                risk_tranches: vec![risk],
924                close_events: Vec::new(),
925                realized_pnl: 0.0,
926                native_realized_pnl: 0.0,
927            },
928        );
929        if let Some(origin) = origin {
930            self.record_pending_terminal(
931                &origin,
932                PendingOrderLifecycleState::Filled,
933                action_id.map(str::to_owned),
934                size,
935                Some(execution.price),
936                fill.ts,
937            );
938        }
939        Ok(())
940    }
941
942    #[allow(clippy::too_many_arguments)]
943    fn record_scale_in(
944        &mut self,
945        id: &str,
946        fill: &FutureFill,
947        engine: &TradeEngine,
948        quote: &PriceQuote,
949        action_id: Option<&str>,
950        signal_ts: Option<NaiveDateTime>,
951        effective_ts: NaiveDateTime,
952        conversion_quotes: Option<&ConversionQuoteBook>,
953    ) -> Result<(), FutureExecutorError> {
954        let position = engine
955            .get_position(id)
956            .ok_or_else(|| FutureExecutorError::PositionNotFound(id.to_owned()))?;
957        if position.data.status != qs_core::types::PositionStatus::Open {
958            return Err(FutureExecutorError::InvalidFill {
959                position_id: id.to_owned(),
960                reason: format!("position is not open: {}", position.data.status),
961            });
962        }
963        if position.data.side != fill.execution.side {
964            return Err(FutureExecutorError::InvalidFill {
965                position_id: id.to_owned(),
966                reason: "execution side does not match position".into(),
967            });
968        }
969        let account = self
970            .accounts
971            .get(id)
972            .ok_or_else(|| FutureExecutorError::AccountNotFound(id.to_owned()))?;
973        let symbol = account.symbol.clone();
974        let side = account.side;
975        let effective_stop = account.effective_stop;
976        let execution = fill.execution;
977        let size = fill.size;
978        let recorded = RecordedFill::from_quote_at(
979            id.to_owned(),
980            action_id.map(str::to_owned),
981            self.fill_sequence,
982            signal_ts,
983            effective_ts,
984            fill.ts,
985            size,
986            quote,
987            execution,
988        );
989        let contract_size = self.contract_size(&symbol);
990        let risk = self.account_risk_tranche(
991            &symbol,
992            fill.ts,
993            RiskTranche::calculate(
994                Some(recorded.id.clone()),
995                side,
996                size,
997                execution.price,
998                effective_stop.map(|stop| stop.price),
999                contract_size,
1000                self.pnl_epsilon,
1001            ),
1002            conversion_quotes,
1003        )?;
1004        self.fill_sequence += 1;
1005        let account = self.accounts.get_mut(id).expect("account checked above");
1006        account.risk_tranches.push(risk);
1007        account.entry_size += size;
1008        account.remaining_size += size;
1009        account.entry_value += execution.price * size;
1010        account.open_entry_value += execution.price * size;
1011        self.fills.push(recorded);
1012        Ok(())
1013    }
1014
1015    #[allow(clippy::too_many_arguments)]
1016    fn record_close(
1017        &mut self,
1018        id: &str,
1019        reason: CloseReason,
1020        fill: &FutureFill,
1021        quote: &PriceQuote,
1022        action_id: Option<&str>,
1023        signal_ts: Option<NaiveDateTime>,
1024        effective_ts: NaiveDateTime,
1025        portfolio_batch: &mut PortfolioBatch,
1026        full_close: bool,
1027        conversion_quotes: Option<&ConversionQuoteBook>,
1028    ) -> Result<(), FutureExecutorError> {
1029        let account = self
1030            .accounts
1031            .get(id)
1032            .ok_or_else(|| FutureExecutorError::AccountNotFound(id.to_owned()))?;
1033        let tolerance = position_size_tolerance(account.entry_size);
1034        if account.remaining_size <= tolerance {
1035            return Err(FutureExecutorError::InvalidFill {
1036                position_id: id.to_owned(),
1037                reason: "position has no remaining size".into(),
1038            });
1039        }
1040        if account.side != fill.execution.side {
1041            return Err(FutureExecutorError::InvalidFill {
1042                position_id: id.to_owned(),
1043                reason: "execution side does not match account".into(),
1044            });
1045        }
1046        if fill.size > account.remaining_size + tolerance {
1047            return Err(FutureExecutorError::InvalidFill {
1048                position_id: id.to_owned(),
1049                reason: format!(
1050                    "close size {} exceeds remaining size {}",
1051                    fill.size, account.remaining_size
1052                ),
1053            });
1054        }
1055        if full_close && (fill.size - account.remaining_size).abs() > tolerance {
1056            return Err(FutureExecutorError::InvalidFill {
1057                position_id: id.to_owned(),
1058                reason: format!(
1059                    "full-close size {} does not consume remaining size {}",
1060                    fill.size, account.remaining_size
1061                ),
1062            });
1063        }
1064
1065        let execution = fill.execution;
1066        let close_size = if full_close {
1067            account.remaining_size
1068        } else {
1069            fill.size.min(account.remaining_size)
1070        };
1071        let next_remaining = (account.remaining_size - close_size).max(0.0);
1072        if !full_close && next_remaining <= tolerance {
1073            return Err(FutureExecutorError::InvalidFill {
1074                position_id: id.to_owned(),
1075                reason: "partial-close effect consumed the entire account".into(),
1076            });
1077        }
1078        let contract_size = self.contract_size(&account.symbol);
1079        let entry_price = account.average_entry();
1080        let native_pnl = match account.side {
1081            Side::Buy => execution.price - entry_price,
1082            Side::Sell => entry_price - execution.price,
1083        } * close_size
1084            * contract_size;
1085        let accounted =
1086            self.convert_native_amount(&account.symbol, native_pnl, fill.ts, conversion_quotes)?;
1087        let pnl = accounted.amount;
1088        let next_balance = self.balance + pnl;
1089        if !next_balance.is_finite() || !portfolio_batch.add_realized_pnl(pnl) {
1090            return Err(FutureExecutorError::PortfolioRejectedRealizedPnl {
1091                position_id: id.to_owned(),
1092                pnl,
1093            });
1094        }
1095
1096        let recorded = RecordedFill::from_quote_at(
1097            id.to_owned(),
1098            action_id.map(str::to_owned),
1099            self.fill_sequence,
1100            signal_ts,
1101            effective_ts,
1102            fill.ts,
1103            close_size,
1104            quote,
1105            execution,
1106        );
1107        self.fill_sequence += 1;
1108
1109        let account = self.accounts.get_mut(id).expect("account checked above");
1110        account.remaining_size = if full_close { 0.0 } else { next_remaining };
1111        account.open_entry_value = if full_close {
1112            0.0
1113        } else {
1114            (account.open_entry_value - entry_price * close_size).max(0.0)
1115        };
1116        account.realized_pnl += pnl;
1117        account.native_realized_pnl += native_pnl;
1118
1119        let mut event = CloseEvent::new(
1120            id.to_owned(),
1121            self.close_sequence,
1122            account.symbol.clone(),
1123            account.side,
1124            fill.ts,
1125            close_size,
1126            execution.price,
1127            pnl,
1128            reason,
1129        );
1130        self.close_sequence += 1;
1131        event.action_id = action_id.map(str::to_owned);
1132        event.fill_id = Some(recorded.id.clone());
1133        event.entry_price = Some(entry_price);
1134        event.native_pnl = Some(native_pnl);
1135        event.native_currency = accounted.native_currency;
1136        event.pnl_conversion = accounted.conversion;
1137        event.remaining_size = Some(account.remaining_size);
1138        account.close_events.push(event.clone());
1139        let trade = TradeResult {
1140            position_id: id.to_owned(),
1141            symbol: account.symbol.clone(),
1142            side: account.side,
1143            entry_price,
1144            exit_price: execution.price,
1145            size: close_size,
1146            pnl,
1147            open_ts: account.open_ts,
1148            close_ts: fill.ts,
1149            close_reason: reason,
1150            group: account.group.clone(),
1151        };
1152
1153        self.balance = next_balance;
1154        self.fills.push(recorded);
1155        self.close_events.push(event);
1156        self.trade_log.push(trade);
1157
1158        if full_close {
1159            let account = self.accounts.remove(id).expect("completed account exists");
1160            let final_net_pnl = account.realized_pnl;
1161            let average_entry = account.historical_average_entry();
1162            let mut completed = CompletedPosition::from_close_events(
1163                account.position_id,
1164                account.symbol,
1165                account.side,
1166                account.open_ts,
1167                fill.ts,
1168                account.entry_size,
1169                average_entry,
1170                account.initial_stop,
1171                account.effective_stop,
1172                account.risk_tranches,
1173                account.close_events,
1174                None,
1175                None,
1176                self.pnl_epsilon,
1177            );
1178            completed.group = account.group;
1179            completed.trade_id = account.trade_id;
1180            let completed_position_index = self.completed_positions.len();
1181            self.completed_positions.push(completed);
1182            portfolio_batch.finish_campaign(id.to_owned(), final_net_pnl, completed_position_index);
1183        }
1184        Ok(())
1185    }
1186
1187    fn account_risk_tranche(
1188        &self,
1189        symbol: &str,
1190        operation_ts: NaiveDateTime,
1191        mut tranche: RiskTranche,
1192        conversion_quotes: Option<&ConversionQuoteBook>,
1193    ) -> Result<RiskTranche, FutureExecutorError> {
1194        let Some(plan) = self.currency_plan.as_ref() else {
1195            return Ok(tranche);
1196        };
1197        let native_currency = plan
1198            .pnl_currency_for_primary_symbol(symbol)
1199            .ok_or_else(|| FutureExecutorError::MissingCurrencyRoute(symbol.to_owned()))?;
1200        tranche.native_currency = Some(native_currency.to_owned());
1201        if tranche.status != RiskBasisStatus::Available {
1202            return Ok(tranche);
1203        }
1204        let Some(native_risk) = tranche.native_risk_amount else {
1205            return Ok(tranche);
1206        };
1207        let accounted =
1208            self.convert_native_amount(symbol, -native_risk, operation_ts, conversion_quotes)?;
1209        let account_risk = -accounted.amount;
1210        if !account_risk.is_finite() || account_risk < 0.0 {
1211            return Err(FutureExecutorError::InvalidConvertedAmount {
1212                symbol: symbol.to_owned(),
1213                kind: "risk",
1214                amount: account_risk,
1215            });
1216        }
1217        tranche.risk_amount = Some(account_risk);
1218        tranche.risk_conversion = accounted.conversion;
1219        Ok(tranche)
1220    }
1221
1222    fn convert_native_amount(
1223        &self,
1224        symbol: &str,
1225        amount: f64,
1226        operation_ts: NaiveDateTime,
1227        conversion_quotes: Option<&ConversionQuoteBook>,
1228    ) -> Result<AccountedAmount, FutureExecutorError> {
1229        let Some(plan) = self.currency_plan.as_ref() else {
1230            return Ok(AccountedAmount {
1231                amount,
1232                native_currency: None,
1233                conversion: None,
1234            });
1235        };
1236        let native_currency = plan
1237            .pnl_currency_for_primary_symbol(symbol)
1238            .ok_or_else(|| FutureExecutorError::MissingCurrencyRoute(symbol.to_owned()))?;
1239        let route = plan
1240            .route_for_primary_symbol(symbol)
1241            .ok_or_else(|| FutureExecutorError::MissingCurrencyRoute(symbol.to_owned()))?;
1242        let conversion = match conversion_quotes {
1243            Some(quotes) => quotes.convert_route(amount, operation_ts, route),
1244            None => match route {
1245                ConversionRoute::Identity { .. } => Ok(ConversionResult {
1246                    from_currency: route.from_currency().to_owned(),
1247                    to_currency: route.to_currency().to_owned(),
1248                    input_amount: amount,
1249                    output_amount: amount,
1250                    operation_ts,
1251                    route: route.clone(),
1252                    legs: Vec::new(),
1253                }),
1254                _ => Err(ConversionError::NoCausalQuote {
1255                    symbol: route.symbols().next().unwrap_or(symbol).to_owned(),
1256                    operation_ts,
1257                    next_quote_ts: None,
1258                }),
1259            },
1260        }
1261        .map_err(|source| FutureExecutorError::Conversion {
1262            symbol: symbol.to_owned(),
1263            operation_ts,
1264            source,
1265        })?;
1266        if !conversion.output_amount.is_finite() {
1267            return Err(FutureExecutorError::InvalidConvertedAmount {
1268                symbol: symbol.to_owned(),
1269                kind: "P&L",
1270                amount: conversion.output_amount,
1271            });
1272        }
1273        Ok(AccountedAmount {
1274            amount: conversion.output_amount,
1275            native_currency: Some(native_currency.to_owned()),
1276            conversion: Some(conversion),
1277        })
1278    }
1279
1280    fn contract_size(&self, symbol: &str) -> f64 {
1281        self.contract_sizes.get(symbol).copied().unwrap_or(1.0)
1282    }
1283}
1284
1285#[cfg(test)]
1286mod tests {
1287    use super::*;
1288    use chrono::{Duration, NaiveDate};
1289    use qs_core::types::{
1290        Action, ExecutionFill, FillPurpose, OrderType, PositionStatus, TargetSpec,
1291    };
1292
1293    use crate::currency::{ConversionPriceSide, FxPair};
1294
1295    fn ts() -> NaiveDateTime {
1296        NaiveDate::from_ymd_opt(2026, 1, 1)
1297            .unwrap()
1298            .and_hms_opt(10, 0, 0)
1299            .unwrap()
1300    }
1301
1302    fn quote_at(seconds: i64, price: f64) -> PriceQuote {
1303        quote_for("EURUSD", seconds, price, price)
1304    }
1305
1306    fn quote_for(symbol: &str, seconds: i64, bid: f64, ask: f64) -> PriceQuote {
1307        PriceQuote {
1308            symbol: symbol.into(),
1309            ts: ts() + Duration::seconds(seconds),
1310            bid,
1311            ask,
1312        }
1313    }
1314
1315    fn eur_account_plan() -> RunCurrencyPlan {
1316        RunCurrencyPlan::new(
1317            "USD",
1318            ["PRIMARY".to_owned()].into_iter().collect(),
1319            ["EURUSD".to_owned()].into_iter().collect(),
1320            [("PRIMARY".to_owned(), "EUR".to_owned())]
1321                .into_iter()
1322                .collect(),
1323            [(
1324                "EUR".to_owned(),
1325                ConversionRoute::Direct {
1326                    pair: FxPair {
1327                        symbol: "EURUSD".to_owned(),
1328                        base_currency: "EUR".to_owned(),
1329                        quote_currency: "USD".to_owned(),
1330                    },
1331                },
1332            )]
1333            .into_iter()
1334            .collect(),
1335            Vec::new(),
1336        )
1337        .unwrap()
1338    }
1339
1340    fn execution(purpose: FillPurpose, price: f64) -> ExecutionFill {
1341        ExecutionFill {
1342            purpose,
1343            side: Side::Buy,
1344            price,
1345            quote_price: price,
1346            requested_price: None,
1347            slippage_pips: 0.0,
1348        }
1349    }
1350
1351    #[test]
1352    fn close_and_initial_risk_use_signed_account_conversion() {
1353        let plan = eur_account_plan();
1354        let mut engine =
1355            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1356        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9)
1357            .with_currency_plan(Some(plan.clone()));
1358        let mut portfolio =
1359            PortfolioRecorder::new(10_000.0, HashMap::new()).with_currency_plan(Some(plan));
1360        let mut conversions = ConversionQuoteBook::new(Duration::hours(1)).unwrap();
1361        conversions
1362            .record_canonical_tick(quote_for("EURUSD", 0, 2.0, 3.0))
1363            .unwrap();
1364        conversions
1365            .record_canonical_tick(quote_for("EURUSD", 1, 2.0, 3.0))
1366            .unwrap();
1367
1368        let open_quote = quote_for("PRIMARY", 0, 100.0, 100.0);
1369        let effects = engine
1370            .apply_priced_future_action(
1371                Action::Open {
1372                    symbol: "PRIMARY".into(),
1373                    side: Side::Buy,
1374                    order_type: OrderType::Market,
1375                    price: None,
1376                    size: 1.0,
1377                    stoploss: Some(90.0),
1378                    targets: vec![],
1379                    rules: vec![],
1380                    group: None,
1381                    trade_id: None,
1382                },
1383                &open_quote,
1384                execution(FillPurpose::MarketEntry, 100.0),
1385            )
1386            .unwrap();
1387        let id = match effects[0].effect() {
1388            Effect::PositionOpened { id } => id.clone(),
1389            effect => panic!("unexpected effect: {effect:?}"),
1390        };
1391        executor
1392            .process_future_effects_with_currency(
1393                &effects,
1394                &engine,
1395                &open_quote,
1396                Some("open"),
1397                Some(open_quote.ts),
1398                open_quote.ts,
1399                &mut portfolio,
1400                Some(&conversions),
1401            )
1402            .unwrap();
1403
1404        let close_quote = quote_for("PRIMARY", 1, 110.0, 110.0);
1405        let effects = engine
1406            .apply_priced_future_action(
1407                Action::ClosePosition { position_id: id },
1408                &close_quote,
1409                execution(FillPurpose::MarketExit, 110.0),
1410            )
1411            .unwrap();
1412        executor
1413            .process_future_effects_with_currency(
1414                &effects,
1415                &engine,
1416                &close_quote,
1417                Some("close"),
1418                Some(close_quote.ts),
1419                close_quote.ts,
1420                &mut portfolio,
1421                Some(&conversions),
1422            )
1423            .unwrap();
1424
1425        assert_eq!(executor.realized_pnl(), 20.0);
1426        assert_eq!(executor.trade_log[0].pnl, 20.0);
1427        let close = &executor.close_events[0];
1428        assert_eq!(close.native_pnl, Some(10.0));
1429        assert_eq!(close.native_currency.as_deref(), Some("EUR"));
1430        let pnl_conversion = close.pnl_conversion.as_ref().unwrap();
1431        assert_eq!(pnl_conversion.input_amount, 10.0);
1432        assert_eq!(pnl_conversion.output_amount, 20.0);
1433        assert_eq!(pnl_conversion.legs[0].price_side, ConversionPriceSide::Bid);
1434
1435        let risk = &executor.completed_positions[0].risk_tranches[0];
1436        assert_eq!(risk.native_risk_amount, Some(10.0));
1437        assert_eq!(risk.risk_amount, Some(30.0));
1438        let risk_conversion = risk.risk_conversion.as_ref().unwrap();
1439        assert_eq!(risk_conversion.input_amount, -10.0);
1440        assert_eq!(risk_conversion.output_amount, -30.0);
1441        assert_eq!(risk_conversion.legs[0].price_side, ConversionPriceSide::Ask);
1442        assert_eq!(executor.completed_positions[0].realized_r, Some(2.0 / 3.0));
1443    }
1444
1445    #[test]
1446    fn stale_close_conversion_commits_no_accounting_artifacts() {
1447        let plan = eur_account_plan();
1448        let mut engine =
1449            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1450        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9)
1451            .with_currency_plan(Some(plan.clone()));
1452        let mut portfolio =
1453            PortfolioRecorder::new(10_000.0, HashMap::new()).with_currency_plan(Some(plan));
1454        let mut conversions = ConversionQuoteBook::new(Duration::zero()).unwrap();
1455        conversions
1456            .record_canonical_tick(quote_for("EURUSD", 0, 2.0, 3.0))
1457            .unwrap();
1458
1459        let open_quote = quote_for("PRIMARY", 0, 100.0, 100.0);
1460        let effects = engine
1461            .apply_priced_future_action(
1462                Action::Open {
1463                    symbol: "PRIMARY".into(),
1464                    side: Side::Buy,
1465                    order_type: OrderType::Market,
1466                    price: None,
1467                    size: 1.0,
1468                    stoploss: None,
1469                    targets: vec![],
1470                    rules: vec![],
1471                    group: None,
1472                    trade_id: None,
1473                },
1474                &open_quote,
1475                execution(FillPurpose::MarketEntry, 100.0),
1476            )
1477            .unwrap();
1478        let id = match effects[0].effect() {
1479            Effect::PositionOpened { id } => id.clone(),
1480            effect => panic!("unexpected effect: {effect:?}"),
1481        };
1482        executor
1483            .process_future_effects_with_currency(
1484                &effects,
1485                &engine,
1486                &open_quote,
1487                None,
1488                None,
1489                open_quote.ts,
1490                &mut portfolio,
1491                Some(&conversions),
1492            )
1493            .unwrap();
1494        portfolio.record_quote(open_quote.clone());
1495        portfolio.record_with_currency(
1496            open_quote.ts,
1497            executor.open_snapshots(),
1498            Some(&conversions),
1499        );
1500        assert!(portfolio.campaign_excursion(&id).is_some());
1501
1502        let close_quote = quote_for("PRIMARY", 1, 110.0, 110.0);
1503        let effects = engine
1504            .apply_priced_future_action(
1505                Action::ClosePosition {
1506                    position_id: id.clone(),
1507                },
1508                &close_quote,
1509                execution(FillPurpose::MarketExit, 110.0),
1510            )
1511            .unwrap();
1512        let error = executor
1513            .process_future_effects_with_currency(
1514                &effects,
1515                &engine,
1516                &close_quote,
1517                None,
1518                None,
1519                close_quote.ts,
1520                &mut portfolio,
1521                Some(&conversions),
1522            )
1523            .unwrap_err();
1524
1525        assert!(matches!(error, FutureExecutorError::Conversion { .. }));
1526        assert_eq!(executor.balance(), 10_000.0);
1527        assert_eq!(executor.fills.len(), 1);
1528        assert!(executor.close_events.is_empty());
1529        assert!(executor.completed_positions.is_empty());
1530        assert_eq!(portfolio.realized_pnl(), 0.0);
1531        assert!(portfolio.campaign_excursion(&id).is_some());
1532    }
1533
1534    #[test]
1535    fn partial_close_scale_in_and_final_close_use_remaining_average_cost() {
1536        let mut engine =
1537            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1538        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9);
1539        let mut portfolio = PortfolioRecorder::new(10_000.0, HashMap::new());
1540
1541        let open_quote = quote_at(0, 100.0);
1542        let effects = engine
1543            .apply_priced_future_action(
1544                Action::Open {
1545                    symbol: "EURUSD".into(),
1546                    side: Side::Buy,
1547                    order_type: OrderType::Market,
1548                    price: None,
1549                    size: 2.0,
1550                    stoploss: None,
1551                    targets: vec![],
1552                    rules: vec![],
1553                    group: None,
1554                    trade_id: None,
1555                },
1556                &open_quote,
1557                execution(FillPurpose::MarketEntry, 100.0),
1558            )
1559            .unwrap();
1560        let id = match effects[0].effect() {
1561            Effect::PositionOpened { id } => id.clone(),
1562            effect => panic!("unexpected effect: {effect:?}"),
1563        };
1564        executor
1565            .process_future_effects(
1566                &effects,
1567                &engine,
1568                &open_quote,
1569                Some("open"),
1570                Some(open_quote.ts),
1571                open_quote.ts,
1572                &mut portfolio,
1573            )
1574            .unwrap();
1575
1576        let partial_quote = quote_at(1, 110.0);
1577        let effects = engine
1578            .apply_priced_future_action(
1579                Action::ClosePartial {
1580                    position_id: id.clone(),
1581                    ratio: 0.5,
1582                },
1583                &partial_quote,
1584                execution(FillPurpose::MarketExit, 110.0),
1585            )
1586            .unwrap();
1587        executor
1588            .process_future_effects(
1589                &effects,
1590                &engine,
1591                &partial_quote,
1592                Some("partial"),
1593                Some(partial_quote.ts),
1594                partial_quote.ts,
1595                &mut portfolio,
1596            )
1597            .unwrap();
1598
1599        let scale_quote = quote_at(2, 120.0);
1600        let effects = engine
1601            .apply_priced_future_action(
1602                Action::ScaleIn {
1603                    position_id: id.clone(),
1604                    price: None,
1605                    size: 1.0,
1606                    trade_id: None,
1607                },
1608                &scale_quote,
1609                execution(FillPurpose::MarketEntry, 120.0),
1610            )
1611            .unwrap();
1612        executor
1613            .process_future_effects(
1614                &effects,
1615                &engine,
1616                &scale_quote,
1617                Some("scale"),
1618                Some(scale_quote.ts),
1619                scale_quote.ts,
1620                &mut portfolio,
1621            )
1622            .unwrap();
1623        assert_eq!(executor.open_snapshots()[0].average_entry_price, 110.0);
1624
1625        let final_quote = quote_at(3, 130.0);
1626        let effects = engine
1627            .apply_priced_future_action(
1628                Action::ClosePosition {
1629                    position_id: id.clone(),
1630                },
1631                &final_quote,
1632                execution(FillPurpose::MarketExit, 130.0),
1633            )
1634            .unwrap();
1635        executor
1636            .process_future_effects(
1637                &effects,
1638                &engine,
1639                &final_quote,
1640                Some("close"),
1641                Some(final_quote.ts),
1642                final_quote.ts,
1643                &mut portfolio,
1644            )
1645            .unwrap();
1646
1647        assert!(executor.open_snapshots().is_empty());
1648        assert_eq!(executor.close_events.len(), 2);
1649        assert_eq!(executor.close_events[0].entry_price, Some(100.0));
1650        assert_eq!(executor.close_events[0].pnl, 10.0);
1651        assert_eq!(executor.close_events[1].entry_price, Some(110.0));
1652        assert_eq!(executor.close_events[1].pnl, 40.0);
1653        assert_eq!(executor.realized_pnl(), 50.0);
1654        assert_eq!(executor.trade_log[1].entry_price, 110.0);
1655    }
1656
1657    #[test]
1658    fn portfolio_rejection_does_not_commit_executor_close_state() {
1659        let mut engine =
1660            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1661        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9);
1662        let mut portfolio = PortfolioRecorder::new(10_000.0, HashMap::new());
1663        let open_quote = quote_at(0, 100.0);
1664        let effects = engine
1665            .apply_priced_future_action(
1666                Action::Open {
1667                    symbol: "EURUSD".into(),
1668                    side: Side::Buy,
1669                    order_type: OrderType::Market,
1670                    price: None,
1671                    size: 1.0,
1672                    stoploss: None,
1673                    targets: vec![],
1674                    rules: vec![],
1675                    group: None,
1676                    trade_id: None,
1677                },
1678                &open_quote,
1679                execution(FillPurpose::MarketEntry, 100.0),
1680            )
1681            .unwrap();
1682        let id = match effects[0].effect() {
1683            Effect::PositionOpened { id } => id.clone(),
1684            effect => panic!("unexpected effect: {effect:?}"),
1685        };
1686        executor
1687            .process_future_effects(
1688                &effects,
1689                &engine,
1690                &open_quote,
1691                None,
1692                None,
1693                open_quote.ts,
1694                &mut portfolio,
1695            )
1696            .unwrap();
1697        assert!(portfolio.set_realized_pnl(f64::MAX));
1698        let balance_before = executor.balance();
1699        let fills_before = executor.fills.len();
1700
1701        let close_quote = quote_at(1, 1.0e308);
1702        let effects = engine
1703            .apply_priced_future_action(
1704                Action::ClosePosition {
1705                    position_id: id.clone(),
1706                },
1707                &close_quote,
1708                execution(FillPurpose::MarketExit, 1.0e308),
1709            )
1710            .unwrap();
1711        let result = executor.process_future_effects(
1712            &effects,
1713            &engine,
1714            &close_quote,
1715            None,
1716            None,
1717            close_quote.ts,
1718            &mut portfolio,
1719        );
1720
1721        assert!(matches!(
1722            result,
1723            Err(FutureExecutorError::PortfolioRejectedRealizedPnl { .. })
1724        ));
1725        assert_eq!(executor.balance(), balance_before);
1726        assert_eq!(executor.fills.len(), fills_before);
1727        assert!(executor.accounts.contains_key(&id));
1728        assert!(executor.close_events.is_empty());
1729    }
1730
1731    #[test]
1732    fn failed_pending_batch_restores_lifecycle_and_deterministic_sequence() {
1733        let mut engine =
1734            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1735        let placement = engine
1736            .apply_future_action(
1737                Action::Open {
1738                    symbol: "EURUSD".into(),
1739                    side: Side::Buy,
1740                    order_type: OrderType::Limit,
1741                    price: Some(99.0),
1742                    size: 1.0,
1743                    stoploss: Some(95.0),
1744                    targets: vec![],
1745                    rules: vec![],
1746                    group: None,
1747                    trade_id: None,
1748                },
1749                ts(),
1750            )
1751            .unwrap();
1752        let id = effect_position_id(placement[0].effect()).to_owned();
1753        let mut batch = placement.clone();
1754        batch.push(FutureEffect::plain(Effect::StoplossModified {
1755            id: id.clone(),
1756            old_price: 95.0,
1757            new_price: f64::NAN,
1758        }));
1759        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9);
1760        let mut portfolio = PortfolioRecorder::new(10_000.0, HashMap::new());
1761        let quote = quote_at(0, 100.0);
1762
1763        let result = executor.process_future_effects(
1764            &batch,
1765            &engine,
1766            &quote,
1767            Some("pending"),
1768            Some(ts()),
1769            ts(),
1770            &mut portfolio,
1771        );
1772
1773        assert!(matches!(
1774            result,
1775            Err(FutureExecutorError::InvalidFill { .. })
1776        ));
1777        assert!(executor.pending_origins.is_empty());
1778        assert!(executor.pending_order_lifecycle.is_empty());
1779        assert_eq!(executor.pending_lifecycle_sequence, 0);
1780        executor
1781            .process_future_effects(
1782                &placement,
1783                &engine,
1784                &quote,
1785                Some("pending"),
1786                Some(ts()),
1787                ts(),
1788                &mut portfolio,
1789            )
1790            .unwrap();
1791        assert_eq!(executor.pending_order_lifecycle[0].sequence, 0);
1792        assert_eq!(
1793            executor.pending_order_lifecycle[0].id,
1794            deterministic_event_id(&id, "pending_placed", 0)
1795        );
1796    }
1797
1798    #[test]
1799    fn failed_close_batch_restores_executor_and_leaves_portfolio_unchanged() {
1800        let mut engine =
1801            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1802        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9);
1803        let mut portfolio = PortfolioRecorder::new(10_000.0, HashMap::new());
1804        let open_quote = quote_at(0, 100.0);
1805        let open_effects = engine
1806            .apply_priced_future_action(
1807                Action::Open {
1808                    symbol: "EURUSD".into(),
1809                    side: Side::Buy,
1810                    order_type: OrderType::Market,
1811                    price: None,
1812                    size: 1.0,
1813                    stoploss: None,
1814                    targets: vec![],
1815                    rules: vec![],
1816                    group: None,
1817                    trade_id: None,
1818                },
1819                &open_quote,
1820                execution(FillPurpose::MarketEntry, 100.0),
1821            )
1822            .unwrap();
1823        let id = effect_position_id(open_effects[0].effect()).to_owned();
1824        executor
1825            .process_future_effects(
1826                &open_effects,
1827                &engine,
1828                &open_quote,
1829                None,
1830                None,
1831                open_quote.ts,
1832                &mut portfolio,
1833            )
1834            .unwrap();
1835        portfolio.record_quote(open_quote.clone());
1836        portfolio.record(open_quote.ts, executor.open_snapshots());
1837        let campaign_before = portfolio.campaign_excursion(&id);
1838
1839        let close_quote = quote_at(1, 110.0);
1840        let close_effects = engine
1841            .apply_priced_future_action(
1842                Action::ClosePosition {
1843                    position_id: id.clone(),
1844                },
1845                &close_quote,
1846                execution(FillPurpose::MarketExit, 110.0),
1847            )
1848            .unwrap();
1849        let mut batch = close_effects.clone();
1850        batch.push(FutureEffect::plain(Effect::StoplossRemoved {
1851            id: "missing".into(),
1852            old_price: 1.0,
1853        }));
1854        let result = executor.process_future_effects(
1855            &batch,
1856            &engine,
1857            &close_quote,
1858            None,
1859            None,
1860            close_quote.ts,
1861            &mut portfolio,
1862        );
1863
1864        assert!(matches!(result, Err(FutureExecutorError::AccountNotFound(id)) if id == "missing"));
1865        assert_eq!(executor.balance(), 10_000.0);
1866        assert_eq!(executor.fills.len(), 1);
1867        assert!(executor.close_events.is_empty());
1868        assert!(executor.completed_positions.is_empty());
1869        assert!(executor.trade_log.is_empty());
1870        assert!(executor.accounts.contains_key(&id));
1871        assert_eq!(executor.fill_sequence, 1);
1872        assert_eq!(executor.close_sequence, 0);
1873        assert_eq!(portfolio.realized_pnl(), 0.0);
1874        assert_eq!(portfolio.campaign_excursion(&id), campaign_before);
1875
1876        executor
1877            .process_future_effects(
1878                &close_effects,
1879                &engine,
1880                &close_quote,
1881                None,
1882                None,
1883                close_quote.ts,
1884                &mut portfolio,
1885            )
1886            .unwrap();
1887        assert_eq!(executor.fills[1].id, deterministic_event_id(&id, "fill", 1));
1888        assert_eq!(
1889            executor.close_events[0].id,
1890            deterministic_event_id(&id, "close", 0)
1891        );
1892        assert_eq!(portfolio.realized_pnl(), 10.0);
1893        assert!(portfolio.campaign_excursion(&id).is_none());
1894    }
1895
1896    #[test]
1897    fn carried_fill_is_consumed_without_repricing_or_engine_synchronization() {
1898        let quote = PriceQuote {
1899            symbol: "EURUSD".into(),
1900            ts: ts(),
1901            bid: 99.0,
1902            ask: 100.0,
1903        };
1904        let execution = ExecutionFill {
1905            purpose: FillPurpose::MarketEntry,
1906            side: Side::Buy,
1907            price: 123.456,
1908            quote_price: 100.0,
1909            requested_price: None,
1910            slippage_pips: 0.0,
1911        };
1912        let mut engine =
1913            TradeEngine::with_fill_model_and_deterministic_ids(qs_core::types::FillModel::BidAsk);
1914        let effects = engine
1915            .apply_priced_future_action(
1916                Action::Open {
1917                    symbol: "EURUSD".into(),
1918                    side: Side::Buy,
1919                    order_type: OrderType::Market,
1920                    price: Some(1.0),
1921                    size: 2.0,
1922                    stoploss: None,
1923                    targets: Vec::<TargetSpec>::new(),
1924                    rules: vec![],
1925                    group: None,
1926                    trade_id: None,
1927                },
1928                &quote,
1929                execution,
1930            )
1931            .unwrap();
1932        let id = match effects[0].effect() {
1933            Effect::PositionOpened { id } => id.clone(),
1934            effect => panic!("unexpected effect: {effect:?}"),
1935        };
1936
1937        let mut executor = FutureExecutor::new(10_000.0, HashMap::new(), 1.0e-9);
1938        let mut portfolio = PortfolioRecorder::new(10_000.0, HashMap::new());
1939        executor
1940            .process_future_effects(
1941                &effects,
1942                &engine,
1943                &quote,
1944                Some("open"),
1945                Some(ts()),
1946                ts(),
1947                &mut portfolio,
1948            )
1949            .unwrap();
1950
1951        assert_eq!(executor.fills.len(), 1);
1952        assert_eq!(executor.fills[0].fill, execution);
1953        assert_eq!(executor.fills[0].size, 2.0);
1954        let position = engine.get_position(&id).unwrap();
1955        assert_eq!(position.data.status, PositionStatus::Open);
1956        assert_eq!(position.data.entries[0].price, execution.price);
1957        assert_eq!(position.data.entries[0].ts, quote.ts);
1958    }
1959}