1use 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 entry_value: f64,
84 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#[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 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 } 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 } 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 "e,
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 "e,
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 "e,
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 "e,
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}