1use std::collections::VecDeque;
19
20use ahash::{AHashMap, AHashSet};
21use nautilus_model::{
22 identifiers::{ClientOrderId, StrategyId, VenueOrderId},
23 types::Quantity,
24};
25use rust_decimal::Decimal;
26
27use crate::{
28 common::{
29 parse::{make_customer_order_ref, make_customer_order_ref_legacy},
30 types::{BetId, CustomerOrderRef, OrderSyncEntry},
31 },
32 stream::{messages::UnmatchedOrder, parse::FillTracker},
33};
34
35#[derive(Clone, Debug)]
36pub(crate) struct PendingReplaceState {
37 pub(crate) total_quantity: Option<Quantity>,
38 awaiting_reconciliation: bool,
39}
40
41#[derive(Debug, Clone, Copy)]
42struct PendingReductionState {
43 client_order_id: ClientOrderId,
44 original_quantity: Quantity,
45 requested_quantity: Quantity,
46 confirmed_quantity: Option<Quantity>,
47}
48
49impl PendingReductionState {
50 fn is_unconfirmed_for(&self, client_order_id: &ClientOrderId) -> bool {
51 self.client_order_id == *client_order_id && self.confirmed_quantity.is_none()
52 }
53
54 fn can_confirm(&self, client_order_id: &ClientOrderId, active_quantity: Quantity) -> bool {
55 self.is_unconfirmed_for(client_order_id)
56 && active_quantity >= self.requested_quantity
57 && active_quantity < self.original_quantity
58 }
59}
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq)]
62pub(crate) enum CustomerOrderRefResolution {
63 Unique(ClientOrderId),
64 Ambiguous,
65}
66
67impl CustomerOrderRefResolution {
68 pub(crate) fn client_order_id(self) -> Option<ClientOrderId> {
69 match self {
70 Self::Unique(client_order_id) => Some(client_order_id),
71 Self::Ambiguous => None,
72 }
73 }
74}
75
76#[derive(Debug, Clone, Default)]
77struct OrderCorrelation {
78 customer_order_refs: AHashSet<CustomerOrderRef>,
79 strategy_id: Option<StrategyId>,
80 venue_order_id: Option<VenueOrderId>,
81 venue_order_ids: AHashSet<BetId>,
82 accepted: bool,
83 terminal_retained: bool,
84}
85
86#[derive(Debug, Clone)]
87enum TerminalRetentionKey {
88 Owned(ClientOrderId),
89 External(BetId),
90}
91
92#[derive(Clone, Debug, Default)]
102pub struct OcmState {
103 pub fill_tracker: FillTracker,
105 pub terminal_orders: AHashSet<BetId>,
107 pub replaced_venue_order_ids: AHashSet<BetId>,
109 pub(crate) customer_order_refs: AHashMap<CustomerOrderRef, CustomerOrderRefResolution>,
110 order_correlations: AHashMap<ClientOrderId, OrderCorrelation>,
111 terminal_order_queue: VecDeque<TerminalRetentionKey>,
112 canceled_replace_bet_ids: AHashSet<BetId>,
113 pending_replace_state: AHashMap<(ClientOrderId, BetId), PendingReplaceState>,
114 pending_reductions: AHashMap<BetId, PendingReductionState>,
115}
116
117impl OcmState {
118 pub(crate) const DEDUP_RETENTION: usize = 10_000;
120
121 pub(crate) fn register_submission(
122 &mut self,
123 client_order_id: ClientOrderId,
124 strategy_id: StrategyId,
125 ) -> Result<(), String> {
126 self.register_order_ref(client_order_id)?;
127 self.register_order_identity(client_order_id, strategy_id);
128 Ok(())
129 }
130
131 pub fn register_customer_order_ref(&mut self, client_order_id: ClientOrderId) {
133 let _ = self.register_order_ref(client_order_id);
134 }
135
136 pub fn register_customer_order_ref_with_legacy(&mut self, client_order_id: ClientOrderId) {
138 let current = make_customer_order_ref(client_order_id.as_str());
139 let legacy = make_customer_order_ref_legacy(client_order_id.as_str());
140 self.upsert_order_correlation(client_order_id, None, false, [current, legacy]);
141 }
142
143 pub fn register_order_identity(
145 &mut self,
146 client_order_id: ClientOrderId,
147 strategy_id: StrategyId,
148 ) {
149 self.order_correlations
150 .entry(client_order_id)
151 .or_default()
152 .strategy_id = Some(strategy_id);
153 }
154
155 pub(crate) fn register_order_ref(
156 &mut self,
157 client_order_id: ClientOrderId,
158 ) -> Result<(), String> {
159 let customer_order_ref = make_customer_order_ref(client_order_id.as_str());
160
161 if self
162 .customer_order_refs
163 .get(&customer_order_ref)
164 .is_some_and(|resolution| {
165 *resolution != CustomerOrderRefResolution::Unique(client_order_id)
166 })
167 {
168 return Err(customer_order_ref);
169 }
170
171 self.upsert_order_correlation(client_order_id, None, false, [customer_order_ref]);
172 Ok(())
173 }
174
175 pub(crate) fn restore_order(
176 &mut self,
177 client_order_id: ClientOrderId,
178 strategy_id: StrategyId,
179 venue_order_id: VenueOrderId,
180 ) {
181 let current = make_customer_order_ref(client_order_id.as_str());
182 let legacy = make_customer_order_ref_legacy(client_order_id.as_str());
183 self.upsert_order_correlation(client_order_id, Some(strategy_id), true, [current, legacy]);
184 self.bind_venue_order_id(&client_order_id, venue_order_id);
185 }
186
187 fn upsert_order_correlation(
188 &mut self,
189 client_order_id: ClientOrderId,
190 strategy_id: Option<StrategyId>,
191 accepted: bool,
192 customer_order_refs: impl IntoIterator<Item = String>,
193 ) {
194 let correlation = self.order_correlations.entry(client_order_id).or_default();
195 if let Some(strategy_id) = strategy_id {
196 correlation.strategy_id = Some(strategy_id);
197 }
198
199 correlation.accepted |= accepted;
200 let mut affected = Vec::new();
201
202 for customer_order_ref in customer_order_refs {
203 if correlation
204 .customer_order_refs
205 .insert(customer_order_ref.clone())
206 {
207 affected.push(customer_order_ref);
208 }
209 }
210
211 for customer_order_ref in affected {
212 self.add_customer_order_ref(client_order_id, customer_order_ref);
213 }
214 }
215
216 fn add_customer_order_ref(
217 &mut self,
218 client_order_id: ClientOrderId,
219 customer_order_ref: String,
220 ) {
221 match self.customer_order_refs.get(&customer_order_ref).copied() {
222 None => {
223 self.customer_order_refs.insert(
224 customer_order_ref,
225 CustomerOrderRefResolution::Unique(client_order_id),
226 );
227 }
228 Some(CustomerOrderRefResolution::Unique(owner)) if owner != client_order_id => {
229 self.customer_order_refs
230 .insert(customer_order_ref, CustomerOrderRefResolution::Ambiguous);
231 }
232 Some(_) => {}
233 }
234 }
235
236 pub fn order_strategy_id(&self, client_order_id: &ClientOrderId) -> Option<StrategyId> {
238 self.order_correlations
239 .get(client_order_id)
240 .and_then(|correlation| correlation.strategy_id)
241 }
242
243 pub(crate) fn bind_venue_order_id(
244 &mut self,
245 client_order_id: &ClientOrderId,
246 venue_order_id: VenueOrderId,
247 ) {
248 let bet_id = venue_order_id.to_string();
249 let correlation = self.order_correlations.entry(*client_order_id).or_default();
250 let newly_correlated = correlation.venue_order_ids.insert(bet_id.clone());
251 correlation.venue_order_id = Some(venue_order_id);
252
253 if self.terminal_orders.contains(&bet_id) {
254 if newly_correlated {
255 self.remove_external_terminal_retention(&bet_id);
256 }
257
258 if !self.has_pending_replace(client_order_id) {
259 self.retain_terminal_identity(*client_order_id);
260 }
261 }
262 }
263
264 pub fn mark_accepted(&mut self, client_order_id: ClientOrderId) -> bool {
269 let correlation = self.order_correlations.entry(client_order_id).or_default();
270 if correlation.accepted {
271 return false;
272 }
273
274 correlation.accepted = true;
275 true
276 }
277
278 pub(crate) fn is_accepted(&self, client_order_id: &ClientOrderId) -> bool {
279 self.order_correlations
280 .get(client_order_id)
281 .is_some_and(|correlation| correlation.accepted)
282 }
283
284 pub(crate) fn claim_acceptance(
285 &mut self,
286 client_order_id: ClientOrderId,
287 venue_order_id: VenueOrderId,
288 ) -> bool {
289 if !self.mark_accepted(client_order_id) {
290 return false;
291 }
292
293 self.bind_venue_order_id(&client_order_id, venue_order_id);
294 true
295 }
296
297 pub(crate) fn remove_order_correlation(&mut self, client_order_id: &ClientOrderId) {
298 self.remove_terminal_retention(client_order_id);
299 let Some(correlation) = self.order_correlations.remove(client_order_id) else {
300 return;
301 };
302 let OrderCorrelation {
303 customer_order_refs: affected,
304 venue_order_ids,
305 ..
306 } = correlation;
307
308 for customer_order_ref in affected {
309 if self
310 .customer_order_refs
311 .get(&customer_order_ref)
312 .is_some_and(|resolution| {
313 *resolution == CustomerOrderRefResolution::Unique(*client_order_id)
314 })
315 {
316 self.customer_order_refs.remove(&customer_order_ref);
317 } else {
318 self.rebuild_customer_order_ref(&customer_order_ref);
319 }
320 }
321
322 self.pending_reductions
323 .retain(|_, pending| pending.client_order_id != *client_order_id);
324 self.pending_replace_state
325 .retain(|(candidate, _), _| candidate != client_order_id);
326
327 for bet_id in venue_order_ids {
328 self.remove_venue_order_state(&bet_id);
329 }
330 }
331
332 pub fn remove_customer_order_refs(&mut self, client_order_id: &ClientOrderId) {
334 self.remove_order_correlation(client_order_id);
335 }
336
337 fn rebuild_customer_order_ref(&mut self, customer_order_ref: &str) {
338 let mut owners = self
339 .order_correlations
340 .iter()
341 .filter(|(_, correlation)| correlation.customer_order_refs.contains(customer_order_ref))
342 .map(|(client_order_id, _)| *client_order_id);
343 let resolution = owners.next().map(|client_order_id| {
344 if owners.next().is_some() {
345 CustomerOrderRefResolution::Ambiguous
346 } else {
347 CustomerOrderRefResolution::Unique(client_order_id)
348 }
349 });
350
351 if let Some(resolution) = resolution {
352 self.customer_order_refs
353 .insert(customer_order_ref.to_string(), resolution);
354 } else {
355 self.customer_order_refs.remove(customer_order_ref);
356 }
357 }
358
359 pub(crate) fn customer_order_ref_resolution(
360 &self,
361 customer_order_ref: &str,
362 ) -> Option<CustomerOrderRefResolution> {
363 self.customer_order_refs.get(customer_order_ref).copied()
364 }
365
366 pub(crate) fn resolve_order_owner(
367 &self,
368 customer_order_ref: Option<&str>,
369 venue_order_id: &str,
370 ) -> Option<CustomerOrderRefResolution> {
371 customer_order_ref
372 .and_then(|reference| self.customer_order_ref_resolution(reference))
373 .or_else(|| {
374 self.client_order_id_by_venue_order_id(venue_order_id)
375 .map(CustomerOrderRefResolution::Unique)
376 })
377 }
378
379 pub(crate) fn client_order_id_by_venue_order_id(
380 &self,
381 venue_order_id: &str,
382 ) -> Option<ClientOrderId> {
383 let mut owners = self
384 .order_correlations
385 .iter()
386 .filter(|(_, correlation)| correlation.venue_order_ids.contains(venue_order_id))
387 .map(|(client_order_id, _)| *client_order_id);
388 let owner = owners.next()?;
389 owners.next().is_none().then_some(owner)
390 }
391
392 pub fn resolve_client_order_id(&self, rfo: Option<&str>) -> Option<ClientOrderId> {
394 rfo.and_then(|customer_order_ref| {
395 self.customer_order_ref_resolution(customer_order_ref)
396 .and_then(CustomerOrderRefResolution::client_order_id)
397 })
398 }
399
400 pub fn should_suppress_cancel(&self, client_order_id: &ClientOrderId, bet_id: &str) -> bool {
403 if self.replaced_venue_order_ids.contains(bet_id) {
404 return true;
405 }
406
407 self.pending_replace_state
408 .contains_key(&(*client_order_id, bet_id.to_string()))
409 }
410
411 pub(crate) fn is_retained_terminal_order(&self, client_order_id: &ClientOrderId) -> bool {
412 self.order_correlations
413 .get(client_order_id)
414 .is_some_and(|correlation| correlation.terminal_retained)
415 }
416
417 pub(crate) fn should_suppress_replaced_report(&self, bet_id: &str) -> bool {
418 if !self.replaced_venue_order_ids.contains(bet_id) {
419 return false;
420 }
421
422 self.client_order_id_by_venue_order_id(bet_id)
423 .is_none_or(|client_order_id| !self.is_retained_terminal_order(&client_order_id))
424 }
425
426 pub(crate) fn register_pending_replace(
427 &mut self,
428 client_order_id: ClientOrderId,
429 old_bet_id: String,
430 total_quantity: Option<Quantity>,
431 ) {
432 self.remove_terminal_retention(&client_order_id);
433 let key = (client_order_id, old_bet_id);
434
435 self.pending_replace_state.insert(
436 key,
437 PendingReplaceState {
438 total_quantity,
439 awaiting_reconciliation: false,
440 },
441 );
442 }
443
444 pub(crate) fn mark_pending_replace_ambiguous(
445 &mut self,
446 client_order_id: ClientOrderId,
447 old_bet_id: &str,
448 ) {
449 let key = (client_order_id, old_bet_id.to_string());
450
451 if let Some(pending) = self.pending_replace_state.get_mut(&key) {
452 pending.awaiting_reconciliation = true;
453 }
454 }
455
456 pub(crate) fn pending_replace_awaits_reconciliation(
457 &self,
458 client_order_id: &ClientOrderId,
459 old_bet_id: &str,
460 ) -> bool {
461 self.pending_replace_state
462 .get(&(*client_order_id, old_bet_id.to_string()))
463 .is_some_and(|pending| pending.awaiting_reconciliation)
464 }
465
466 pub(crate) fn take_pending_replace(
467 &mut self,
468 client_order_id: ClientOrderId,
469 old_bet_id: &str,
470 ) -> Option<PendingReplaceState> {
471 let key = (client_order_id, old_bet_id.to_string());
472 self.pending_replace_state.remove(&key)
473 }
474
475 fn has_pending_replace(&self, client_order_id: &ClientOrderId) -> bool {
476 self.pending_replace_state
477 .keys()
478 .any(|(candidate, _)| candidate == client_order_id)
479 }
480
481 pub(crate) fn mark_canceled_replace(&mut self, client_order_id: ClientOrderId, bet_id: &str) {
482 self.canceled_replace_bet_ids.insert(bet_id.to_string());
483 self.retain_terminal_order(client_order_id, bet_id);
484 }
485
486 pub(crate) fn is_canceled_replace(&self, bet_id: &str) -> bool {
487 self.canceled_replace_bet_ids.contains(bet_id)
488 }
489
490 pub(crate) fn is_redundant_terminal_update(&self, order: &UnmatchedOrder) -> bool {
491 self.terminal_orders.contains(&order.id)
492 && !self.is_canceled_replace(&order.id)
493 && !self.fill_tracker.has_unseen_fill(order)
494 && !self.fill_tracker.has_unseen_fill_void(order)
495 }
496
497 pub(crate) fn clear_canceled_replace(&mut self, bet_id: &str) {
498 self.canceled_replace_bet_ids.remove(bet_id);
499 }
500
501 pub(crate) fn promote_pending_replace(
506 &mut self,
507 client_order_id: &ClientOrderId,
508 new_bet_id: &str,
509 replacement_quantity: Quantity,
510 ) -> Option<(Quantity, String)> {
511 let mut old_bet_ids = self
512 .pending_replace_state
513 .keys()
514 .filter(|(candidate, old_bet_id)| {
515 candidate == client_order_id && old_bet_id != new_bet_id
516 })
517 .map(|(_, old_bet_id)| old_bet_id.clone());
518
519 let old_bet_id = old_bet_ids.next()?;
520 if old_bet_ids.next().is_some() {
521 return None;
522 }
523
524 let pending = self.complete_pending_replace(
525 *client_order_id,
526 &old_bet_id,
527 VenueOrderId::from(new_bet_id),
528 )?;
529
530 Some((
531 pending.total_quantity.unwrap_or(replacement_quantity),
532 old_bet_id,
533 ))
534 }
535
536 pub(crate) fn complete_pending_replace(
537 &mut self,
538 client_order_id: ClientOrderId,
539 old_bet_id: &str,
540 new_venue_order_id: VenueOrderId,
541 ) -> Option<PendingReplaceState> {
542 if self
543 .replaced_venue_order_ids
544 .contains(new_venue_order_id.as_str())
545 {
546 return None;
547 }
548
549 let pending = self.take_pending_replace(client_order_id, old_bet_id)?;
550 self.mark_replaced_venue_order_id(client_order_id, old_bet_id.to_string());
551 self.bind_venue_order_id(&client_order_id, new_venue_order_id);
552 Some(pending)
553 }
554
555 fn mark_replaced_venue_order_id(&mut self, client_order_id: ClientOrderId, bet_id: String) {
556 self.mark_correlated_terminal_order(client_order_id, &bet_id);
557 self.replaced_venue_order_ids.insert(bet_id);
558 }
559
560 pub(crate) fn register_pending_reduction(
561 &mut self,
562 client_order_id: ClientOrderId,
563 bet_id: String,
564 original_quantity: Quantity,
565 requested_quantity: Quantity,
566 ) {
567 self.pending_reductions.insert(
568 bet_id,
569 PendingReductionState {
570 client_order_id,
571 original_quantity,
572 requested_quantity,
573 confirmed_quantity: None,
574 },
575 );
576 }
577
578 pub(crate) fn confirm_pending_reduction(
579 &mut self,
580 client_order_id: &ClientOrderId,
581 bet_id: &str,
582 active_quantity: Quantity,
583 ) -> Option<Quantity> {
584 let pending = self.pending_reductions.get_mut(bet_id)?;
585
586 if !pending.can_confirm(client_order_id, active_quantity) {
587 return None;
588 }
589
590 pending.confirmed_quantity = Some(active_quantity);
591 Some(active_quantity)
592 }
593
594 pub(crate) fn complete_pending_reduction(
595 &mut self,
596 client_order_id: &ClientOrderId,
597 bet_id: &str,
598 quantity: Quantity,
599 ) -> bool {
600 let Some(pending) = self.pending_reductions.get_mut(bet_id) else {
601 return false;
602 };
603
604 if !pending.is_unconfirmed_for(client_order_id) {
605 return false;
606 }
607
608 pending.confirmed_quantity = Some(quantity);
609 true
610 }
611
612 pub(crate) fn reduced_quantity(&self, bet_id: &str) -> Option<Quantity> {
613 self.pending_reductions
614 .get(bet_id)
615 .and_then(|pending| pending.confirmed_quantity)
616 }
617
618 pub(crate) fn clear_pending_reduction(
619 &mut self,
620 client_order_id: &ClientOrderId,
621 bet_id: &str,
622 ) {
623 if self
624 .pending_reductions
625 .get(bet_id)
626 .is_some_and(|pending| pending.client_order_id == *client_order_id)
627 {
628 self.pending_reductions.remove(bet_id);
629 }
630 }
631
632 pub(crate) fn retain_terminal_order(&mut self, client_order_id: ClientOrderId, bet_id: &str) {
634 if self.replaced_venue_order_ids.contains(bet_id)
635 || self.has_pending_replace(&client_order_id)
636 {
637 self.mark_correlated_terminal_order(client_order_id, bet_id);
638 return;
639 }
640
641 self.bind_venue_order_id(&client_order_id, VenueOrderId::from(bet_id));
642 self.mark_correlated_terminal_order(client_order_id, bet_id);
643
644 self.retain_terminal_identity(client_order_id);
645 }
646
647 fn mark_correlated_terminal_order(&mut self, client_order_id: ClientOrderId, bet_id: &str) {
648 if self
649 .pending_reductions
650 .get(bet_id)
651 .is_some_and(|pending| pending.is_unconfirmed_for(&client_order_id))
652 {
653 self.pending_reductions.remove(bet_id);
654 }
655
656 let newly_correlated = self
657 .order_correlations
658 .entry(client_order_id)
659 .or_default()
660 .venue_order_ids
661 .insert(bet_id.to_string());
662 if newly_correlated {
663 self.remove_external_terminal_retention(bet_id);
664 }
665 self.terminal_orders.insert(bet_id.to_string());
666 }
667
668 fn retain_terminal_identity(&mut self, client_order_id: ClientOrderId) {
669 let correlation = self.order_correlations.entry(client_order_id).or_default();
670 if correlation.terminal_retained {
671 return;
672 }
673
674 correlation.terminal_retained = true;
675 self.push_terminal_identity(TerminalRetentionKey::Owned(client_order_id));
676 }
677
678 pub fn mark_terminal_order(&mut self, bet_id: String) {
680 if !self.terminal_orders.insert(bet_id.clone()) {
681 return;
682 }
683
684 self.push_terminal_identity(TerminalRetentionKey::External(bet_id));
685 }
686
687 fn push_terminal_identity(&mut self, key: TerminalRetentionKey) {
688 self.terminal_order_queue.push_back(key);
689 if self.terminal_order_queue.len() > Self::DEDUP_RETENTION
690 && let Some(expired) = self.terminal_order_queue.pop_front()
691 {
692 self.evict_terminal_identity(expired);
693 }
694 }
695
696 fn remove_terminal_retention(&mut self, client_order_id: &ClientOrderId) {
697 if self
698 .order_correlations
699 .get_mut(client_order_id)
700 .is_some_and(|correlation| std::mem::take(&mut correlation.terminal_retained))
701 {
702 self.terminal_order_queue.retain(
703 |key| !matches!(key, TerminalRetentionKey::Owned(candidate) if candidate == client_order_id),
704 );
705 }
706 }
707
708 fn remove_external_terminal_retention(&mut self, bet_id: &str) {
709 if !self.terminal_orders.contains(bet_id) {
710 return;
711 }
712
713 self.terminal_order_queue.retain(
714 |key| !matches!(key, TerminalRetentionKey::External(candidate) if candidate == bet_id),
715 );
716 }
717
718 pub(crate) fn mark_order_active(&mut self, client_order_id: &ClientOrderId, bet_id: &str) {
719 self.remove_terminal_retention(client_order_id);
720 self.terminal_orders.remove(bet_id);
721 self.canceled_replace_bet_ids.remove(bet_id);
722 }
723
724 fn evict_terminal_identity(&mut self, key: TerminalRetentionKey) {
725 match key {
726 TerminalRetentionKey::Owned(client_order_id) => {
727 let retained = self
728 .order_correlations
729 .get_mut(&client_order_id)
730 .is_some_and(|correlation| std::mem::take(&mut correlation.terminal_retained));
731 if retained {
732 self.remove_order_correlation(&client_order_id);
733 }
734 }
735 TerminalRetentionKey::External(bet_id) => {
736 self.remove_venue_order_state(&bet_id);
737 }
738 }
739 }
740
741 fn remove_venue_order_state(&mut self, bet_id: &str) {
742 self.terminal_orders.remove(bet_id);
743 self.replaced_venue_order_ids.remove(bet_id);
744 self.canceled_replace_bet_ids.remove(bet_id);
745 self.pending_reductions.remove(bet_id);
746 self.fill_tracker.prune(bet_id);
747 }
748
749 pub fn sync_from_orders(&mut self, orders: &[OrderSyncEntry]) {
753 for entry in orders {
754 self.restore_order(
755 entry.client_order_id,
756 entry.strategy_id,
757 VenueOrderId::from(entry.bet_id.as_str()),
758 );
759
760 for venue_order_id in &entry.venue_order_ids {
761 if venue_order_id != &entry.bet_id {
762 self.mark_replaced_venue_order_id(
763 entry.client_order_id,
764 venue_order_id.clone(),
765 );
766 }
767 }
768
769 if entry.is_closed {
770 self.retain_terminal_order(entry.client_order_id, &entry.bet_id);
771 } else {
772 self.mark_order_active(&entry.client_order_id, &entry.bet_id);
773 }
774
775 if entry.filled_qty > Decimal::ZERO {
776 self.fill_tracker
777 .sync_order(&entry.bet_id, entry.filled_qty, entry.avg_px);
778 }
779
780 if !entry.trade_ids.is_empty() {
781 self.fill_tracker
782 .seed_published_trade_ids(entry.trade_ids.iter().cloned());
783 }
784 }
785 }
786}
787
788#[cfg(test)]
789mod tests {
790 use rstest::rstest;
791
792 use super::*;
793
794 #[rstest]
795 fn owned_terminal_retention_evicts_all_identity_state_together() {
796 let mut state = OcmState::default();
797 let client_order_id = ClientOrderId::from("O-TERMINAL-0");
798 let strategy_id = StrategyId::from("S-001");
799 let current_bet_id = "current-bet-0";
800 let replaced_bet_id = "replaced-bet-0";
801 let size_matched = Decimal::from(2);
802 let average_price = Decimal::from(3);
803 state.restore_order(
804 client_order_id,
805 strategy_id,
806 VenueOrderId::from(current_bet_id),
807 );
808 state.mark_replaced_venue_order_id(client_order_id, replaced_bet_id.to_string());
809 assert!(state.is_accepted(&client_order_id));
810 state.register_pending_reduction(
811 client_order_id,
812 current_bet_id.to_string(),
813 Quantity::from(10),
814 Quantity::from(4),
815 );
816 state.confirm_pending_reduction(&client_order_id, current_bet_id, Quantity::from(4));
817 state.fill_tracker.advance_cumulative_fill(
818 current_bet_id,
819 size_matched,
820 Some(average_price),
821 average_price,
822 );
823 state.fill_tracker.advance_cumulative_fill(
824 replaced_bet_id,
825 size_matched,
826 Some(average_price),
827 average_price,
828 );
829 state.retain_terminal_order(client_order_id, current_bet_id);
830
831 for index in 0..OcmState::DEDUP_RETENTION {
832 state.mark_terminal_order(format!("external-bet-{index}"));
833 }
834
835 let current_replay = state.fill_tracker.advance_cumulative_fill(
836 current_bet_id,
837 size_matched,
838 Some(average_price),
839 average_price,
840 );
841 let replaced_replay = state.fill_tracker.advance_cumulative_fill(
842 replaced_bet_id,
843 size_matched,
844 Some(average_price),
845 average_price,
846 );
847 let customer_order_ref = make_customer_order_ref(client_order_id.as_str());
848
849 assert_eq!(state.order_strategy_id(&client_order_id), None);
850 assert_eq!(
851 state.resolve_client_order_id(Some(&customer_order_ref)),
852 None,
853 );
854 assert_eq!(state.reduced_quantity(current_bet_id), None);
855 assert!(!state.terminal_orders.contains(current_bet_id));
856 assert!(!state.terminal_orders.contains(replaced_bet_id));
857 assert!(!state.replaced_venue_order_ids.contains(replaced_bet_id));
858 assert!(current_replay.is_some());
859 assert!(replaced_replay.is_some());
860 }
861
862 #[rstest]
863 fn terminal_retention_does_not_evict_active_or_ambiguous_replace_identity() {
864 let mut state = OcmState::default();
865 let active_client_order_id = ClientOrderId::from("O-ACTIVE");
866 let ambiguous_client_order_id = ClientOrderId::from("O-AMBIGUOUS");
867 let strategy_id = StrategyId::from("S-001");
868 state.restore_order(
869 active_client_order_id,
870 strategy_id,
871 VenueOrderId::from("active-bet"),
872 );
873 state.restore_order(
874 ambiguous_client_order_id,
875 strategy_id,
876 VenueOrderId::from("ambiguous-bet"),
877 );
878 state.register_pending_replace(
879 ambiguous_client_order_id,
880 "ambiguous-bet".to_string(),
881 Some(Quantity::from(10)),
882 );
883 state.mark_pending_replace_ambiguous(ambiguous_client_order_id, "ambiguous-bet");
884 state.retain_terminal_order(ambiguous_client_order_id, "ambiguous-bet");
885
886 for index in 0..=OcmState::DEDUP_RETENTION {
887 state.mark_terminal_order(format!("external-bet-{index}"));
888 }
889
890 assert_eq!(
891 state.order_strategy_id(&active_client_order_id),
892 Some(strategy_id),
893 );
894 assert!(!state.mark_accepted(active_client_order_id));
895 assert_eq!(
896 state.order_strategy_id(&ambiguous_client_order_id),
897 Some(strategy_id),
898 );
899 assert!(state.pending_replace_awaits_reconciliation(
900 &ambiguous_client_order_id,
901 "ambiguous-bet",
902 ));
903 assert!(state.should_suppress_cancel(&ambiguous_client_order_id, "ambiguous-bet",));
904 }
905
906 #[rstest]
907 fn successful_replace_without_old_leg_ocm_uses_terminal_identity_lifecycle() {
908 let mut state = OcmState::default();
909 let client_order_id = ClientOrderId::from("O-REPLACE");
910 let strategy_id = StrategyId::from("S-001");
911 state.restore_order(client_order_id, strategy_id, VenueOrderId::from("old-bet"));
912 state.register_pending_replace(
913 client_order_id,
914 "old-bet".to_string(),
915 Some(Quantity::from(10)),
916 );
917
918 let completed = state.complete_pending_replace(
919 client_order_id,
920 "old-bet",
921 VenueOrderId::from("new-bet"),
922 );
923
924 assert!(completed.is_some());
925 assert!(state.pending_replace_state.is_empty());
926 assert!(state.replaced_venue_order_ids.contains("old-bet"));
927 assert!(state.terminal_orders.contains("old-bet"));
928 assert_eq!(
929 state.client_order_id_by_venue_order_id("old-bet"),
930 Some(client_order_id),
931 );
932 assert_eq!(
933 state.client_order_id_by_venue_order_id("new-bet"),
934 Some(client_order_id),
935 );
936 assert!(!state.is_retained_terminal_order(&client_order_id));
937
938 state.retain_terminal_order(client_order_id, "old-bet");
939
940 assert_eq!(
941 state
942 .order_correlations
943 .get(&client_order_id)
944 .and_then(|correlation| correlation.venue_order_id),
945 Some(VenueOrderId::from("new-bet")),
946 );
947 assert!(!state.is_retained_terminal_order(&client_order_id));
948
949 state.retain_terminal_order(client_order_id, "new-bet");
950 state.retain_terminal_order(client_order_id, "old-bet");
951
952 assert_eq!(
953 state
954 .order_correlations
955 .get(&client_order_id)
956 .and_then(|correlation| correlation.venue_order_id),
957 Some(VenueOrderId::from("new-bet")),
958 );
959 assert_eq!(state.terminal_order_queue.len(), 1);
960 }
961
962 #[rstest]
963 fn replace_resolution_does_not_discard_another_pending_mutation() {
964 let mut state = OcmState::default();
965 let client_order_id = ClientOrderId::from("O-REPLACE");
966 state.restore_order(
967 client_order_id,
968 StrategyId::from("S-001"),
969 VenueOrderId::from("old-bet-1"),
970 );
971 state.register_pending_replace(
972 client_order_id,
973 "old-bet-1".to_string(),
974 Some(Quantity::from(10)),
975 );
976 state.register_pending_replace(
977 client_order_id,
978 "old-bet-2".to_string(),
979 Some(Quantity::from(10)),
980 );
981
982 let unresolved_promotion =
983 state.promote_pending_replace(&client_order_id, "new-bet", Quantity::from(10));
984 let completed = state.complete_pending_replace(
985 client_order_id,
986 "old-bet-1",
987 VenueOrderId::from("new-bet"),
988 );
989
990 assert_eq!(unresolved_promotion, None);
991 assert!(completed.is_some());
992 assert!(
993 !state
994 .pending_replace_state
995 .contains_key(&(client_order_id, "old-bet-1".to_string()))
996 );
997 assert!(
998 state
999 .pending_replace_state
1000 .contains_key(&(client_order_id, "old-bet-2".to_string()))
1001 );
1002 assert!(state.replaced_venue_order_ids.contains("old-bet-1"));
1003 assert!(!state.replaced_venue_order_ids.contains("old-bet-2"));
1004 assert_eq!(
1005 state.client_order_id_by_venue_order_id("new-bet"),
1006 Some(client_order_id),
1007 );
1008 }
1009
1010 #[rstest]
1011 fn historical_bet_migrates_from_external_to_owned_identity() {
1012 let mut state = OcmState::default();
1013 let client_order_id = ClientOrderId::from("O-MIGRATED");
1014 let old_bet_id = "old-bet";
1015 let size_matched = Decimal::from(2);
1016 let average_price = Decimal::from(3);
1017 state.mark_terminal_order(old_bet_id.to_string());
1018 state.fill_tracker.advance_cumulative_fill(
1019 old_bet_id,
1020 size_matched,
1021 Some(average_price),
1022 average_price,
1023 );
1024 state.restore_order(
1025 client_order_id,
1026 StrategyId::from("S-001"),
1027 VenueOrderId::from("current-bet"),
1028 );
1029 state.mark_replaced_venue_order_id(client_order_id, old_bet_id.to_string());
1030
1031 for index in 0..=OcmState::DEDUP_RETENTION {
1032 state.mark_terminal_order(format!("external-bet-{index}"));
1033 }
1034
1035 let replay = state.fill_tracker.advance_cumulative_fill(
1036 old_bet_id,
1037 size_matched,
1038 Some(average_price),
1039 average_price,
1040 );
1041
1042 assert_eq!(
1043 state.client_order_id_by_venue_order_id(old_bet_id),
1044 Some(client_order_id),
1045 );
1046 assert!(state.terminal_orders.contains(old_bet_id));
1047 assert!(state.replaced_venue_order_ids.contains(old_bet_id));
1048 assert!(replay.is_none());
1049 }
1050
1051 #[rstest]
1052 fn sustained_owned_terminal_history_stays_bounded() {
1053 let mut state = OcmState::default();
1054 let strategy_id = StrategyId::from("S-001");
1055
1056 for index in 0..=OcmState::DEDUP_RETENTION {
1057 let client_order_id = ClientOrderId::from(format!("O-{index}"));
1058 let current_bet_id = format!("current-bet-{index}");
1059 state.restore_order(
1060 client_order_id,
1061 strategy_id,
1062 VenueOrderId::from(current_bet_id.as_str()),
1063 );
1064 state.mark_replaced_venue_order_id(client_order_id, format!("replaced-bet-{index}"));
1065 state.retain_terminal_order(client_order_id, ¤t_bet_id);
1066 }
1067
1068 assert_eq!(state.order_correlations.len(), OcmState::DEDUP_RETENTION);
1069 assert_eq!(
1070 state
1071 .order_correlations
1072 .values()
1073 .filter(|correlation| correlation.terminal_retained)
1074 .count(),
1075 OcmState::DEDUP_RETENTION,
1076 );
1077 assert_eq!(state.terminal_order_queue.len(), OcmState::DEDUP_RETENTION);
1078 assert_eq!(
1079 state.replaced_venue_order_ids.len(),
1080 OcmState::DEDUP_RETENTION,
1081 );
1082 assert_eq!(state.terminal_orders.len(), OcmState::DEDUP_RETENTION * 2);
1083 assert_eq!(state.order_strategy_id(&ClientOrderId::from("O-0")), None,);
1084 assert_eq!(
1085 state.order_strategy_id(&ClientOrderId::from(format!(
1086 "O-{}",
1087 OcmState::DEDUP_RETENTION
1088 ))),
1089 Some(strategy_id),
1090 );
1091 }
1092
1093 #[rstest]
1094 fn terminal_order_retention_evicts_fill_tracker_state() {
1095 let mut state = OcmState::default();
1096 let first_bet_id = "bet-0";
1097 let size_matched = Decimal::new(10, 0);
1098 let average_price = Decimal::new(20, 1);
1099
1100 assert!(
1101 state
1102 .fill_tracker
1103 .advance_cumulative_fill(
1104 first_bet_id,
1105 size_matched,
1106 Some(average_price),
1107 average_price,
1108 )
1109 .is_some(),
1110 );
1111 state
1112 .replaced_venue_order_ids
1113 .insert(first_bet_id.to_string());
1114 state.mark_terminal_order(first_bet_id.to_string());
1115 for index in 1..=OcmState::DEDUP_RETENTION {
1116 state.mark_terminal_order(format!("bet-{index}"));
1117 }
1118
1119 let replay_after_eviction = state.fill_tracker.advance_cumulative_fill(
1120 first_bet_id,
1121 size_matched,
1122 Some(average_price),
1123 average_price,
1124 );
1125
1126 assert!(!state.terminal_orders.contains(first_bet_id));
1127 assert!(!state.replaced_venue_order_ids.contains(first_bet_id));
1128 assert!(replay_after_eviction.is_some());
1129 }
1130
1131 #[rstest]
1132 fn active_submission_collision_is_rejected() {
1133 let suffix = "12345678901234567890123456789012";
1134 let first = ClientOrderId::from(format!("FIRST-{suffix}"));
1135 let second = ClientOrderId::from(format!("SECOND-{suffix}"));
1136 let strategy_id = StrategyId::from("S-001");
1137 let mut state = OcmState::default();
1138
1139 state.register_submission(first, strategy_id).unwrap();
1140 let collision = state.register_submission(second, strategy_id);
1141
1142 assert_eq!(collision, Err(suffix.to_string()));
1143 assert_eq!(state.resolve_client_order_id(Some(suffix)), Some(first));
1144 assert_eq!(state.order_strategy_id(&second), None);
1145 }
1146
1147 #[rstest]
1148 fn submission_collision_with_restored_legacy_reference_is_rejected() {
1149 let reference = "12345678901234567890123456789012";
1150 let restored = ClientOrderId::from(format!("{reference}-RESTORED"));
1151 let fresh = ClientOrderId::from(format!("FRESH-{reference}"));
1152 let strategy_id = StrategyId::from("S-001");
1153 let mut state = OcmState::default();
1154
1155 state.restore_order(restored, strategy_id, VenueOrderId::from("bet-restored"));
1156 let collision = state.register_submission(fresh, strategy_id);
1157
1158 assert_eq!(collision, Err(reference.to_string()));
1159 assert_eq!(
1160 state.resolve_client_order_id(Some(reference)),
1161 Some(restored),
1162 );
1163 assert_eq!(state.order_strategy_id(&fresh), None);
1164 }
1165
1166 #[rstest]
1167 fn restored_current_reference_ambiguity_recovers_after_cleanup() {
1168 let suffix = "12345678901234567890123456789012";
1169 let first = ClientOrderId::from(format!("FIRST-{suffix}"));
1170 let second = ClientOrderId::from(format!("SECOND-{suffix}"));
1171 let strategy_id = StrategyId::from("S-001");
1172 let mut state = OcmState::default();
1173
1174 state.restore_order(first, strategy_id, VenueOrderId::from("bet-1"));
1175 state.restore_order(second, strategy_id, VenueOrderId::from("bet-2"));
1176
1177 assert_eq!(
1178 state.customer_order_ref_resolution(suffix),
1179 Some(CustomerOrderRefResolution::Ambiguous),
1180 );
1181 assert_eq!(state.resolve_client_order_id(Some(suffix)), None);
1182 state.remove_order_correlation(&first);
1183 assert_eq!(state.resolve_client_order_id(Some(suffix)), Some(second));
1184 assert_eq!(
1185 state
1186 .order_correlations
1187 .get(&second)
1188 .and_then(|correlation| correlation.venue_order_id),
1189 Some(VenueOrderId::from("bet-2")),
1190 );
1191 assert!(!state.mark_accepted(second));
1192 }
1193
1194 #[rstest]
1195 fn restored_legacy_reference_ambiguity_recovers_after_cleanup() {
1196 let prefix = "12345678901234567890123456789012";
1197 let first = ClientOrderId::from(format!("{prefix}-FIRST"));
1198 let second = ClientOrderId::from(format!("{prefix}-SECOND"));
1199 let strategy_id = StrategyId::from("S-001");
1200 let mut state = OcmState::default();
1201
1202 state.restore_order(first, strategy_id, VenueOrderId::from("bet-1"));
1203 state.restore_order(second, strategy_id, VenueOrderId::from("bet-2"));
1204
1205 assert_eq!(
1206 state.customer_order_ref_resolution(prefix),
1207 Some(CustomerOrderRefResolution::Ambiguous),
1208 );
1209 assert_eq!(state.resolve_client_order_id(Some(prefix)), None);
1210 state.remove_order_correlation(&second);
1211 assert_eq!(state.resolve_client_order_id(Some(prefix)), Some(first));
1212 }
1213
1214 #[rstest]
1215 fn venue_identity_can_be_bound_and_replaced() {
1216 let client_order_id = ClientOrderId::from("O-1");
1217 let mut state = OcmState::default();
1218 state
1219 .register_submission(client_order_id, StrategyId::from("S-001"))
1220 .unwrap();
1221
1222 state.bind_venue_order_id(&client_order_id, VenueOrderId::from("bet-1"));
1223 assert_eq!(
1224 state
1225 .order_correlations
1226 .get(&client_order_id)
1227 .and_then(|correlation| correlation.venue_order_id),
1228 Some(VenueOrderId::from("bet-1")),
1229 );
1230
1231 state.bind_venue_order_id(&client_order_id, VenueOrderId::from("bet-2"));
1232 assert_eq!(
1233 state
1234 .order_correlations
1235 .get(&client_order_id)
1236 .and_then(|correlation| correlation.venue_order_id),
1237 Some(VenueOrderId::from("bet-2")),
1238 );
1239 }
1240
1241 #[rstest]
1242 fn claimed_acceptance_does_not_replace_restored_venue_identity() {
1243 let client_order_id = ClientOrderId::from("O-1");
1244 let mut state = OcmState::default();
1245 state.restore_order(
1246 client_order_id,
1247 StrategyId::from("S-001"),
1248 VenueOrderId::from("bet-current"),
1249 );
1250
1251 let claimed = state.claim_acceptance(client_order_id, VenueOrderId::from("bet-stale"));
1252
1253 assert!(!claimed);
1254 assert_eq!(
1255 state
1256 .order_correlations
1257 .get(&client_order_id)
1258 .and_then(|correlation| correlation.venue_order_id),
1259 Some(VenueOrderId::from("bet-current")),
1260 );
1261 }
1262
1263 #[rstest]
1264 fn pending_replace_promotes_only_a_different_bet() {
1265 let client_order_id = ClientOrderId::from("O-1");
1266 let mut state = OcmState::default();
1267 state.restore_order(
1268 client_order_id,
1269 StrategyId::from("S-001"),
1270 VenueOrderId::from("old-bet"),
1271 );
1272 state.register_pending_replace(
1273 client_order_id,
1274 "old-bet".to_string(),
1275 Some(Quantity::from(10)),
1276 );
1277
1278 assert_eq!(
1279 state.promote_pending_replace(&client_order_id, "old-bet", Quantity::from(8)),
1280 None
1281 );
1282 assert_eq!(
1283 state.promote_pending_replace(&client_order_id, "new-bet", Quantity::from(8)),
1284 Some((Quantity::from(10), "old-bet".to_string())),
1285 );
1286 assert!(state.pending_replace_state.is_empty());
1287 assert!(state.replaced_venue_order_ids.contains("old-bet"));
1288 assert_eq!(
1289 state
1290 .order_correlations
1291 .get(&client_order_id)
1292 .and_then(|correlation| correlation.venue_order_id),
1293 Some(VenueOrderId::from("new-bet")),
1294 );
1295 assert_eq!(
1296 state.promote_pending_replace(&client_order_id, "newer-bet", Quantity::from(8)),
1297 None,
1298 );
1299 }
1300
1301 #[rstest]
1302 fn pending_replace_does_not_promote_a_historical_bet() {
1303 let client_order_id = ClientOrderId::from("O-1");
1304 let mut state = OcmState::default();
1305 state
1306 .replaced_venue_order_ids
1307 .insert("historical-bet".to_string());
1308 state.register_pending_replace(
1309 client_order_id,
1310 "current-bet".to_string(),
1311 Some(Quantity::from(10)),
1312 );
1313
1314 assert_eq!(
1315 state.promote_pending_replace(&client_order_id, "historical-bet", Quantity::from(8)),
1316 None,
1317 );
1318 assert!(
1319 state.should_suppress_cancel(&client_order_id, "current-bet"),
1320 "historical traffic must not consume the pending replace",
1321 );
1322 assert_eq!(
1323 state.promote_pending_replace(&client_order_id, "replacement-bet", Quantity::from(8)),
1324 Some((Quantity::from(10), "current-bet".to_string())),
1325 );
1326 }
1327
1328 #[rstest]
1329 #[case::below_requested(Quantity::from(3), None)]
1330 #[case::at_requested(Quantity::from(4), Some(Quantity::from(4)))]
1331 #[case::inside_window(Quantity::from(9), Some(Quantity::from(9)))]
1332 #[case::at_original(Quantity::from(10), None)]
1333 fn pending_reduction_confirms_only_a_definitive_reduction(
1334 #[case] active_quantity: Quantity,
1335 #[case] expected: Option<Quantity>,
1336 ) {
1337 let client_order_id = ClientOrderId::from("O-1");
1338 let bet_id = "bet-1";
1339 let mut state = OcmState::default();
1340 state.register_pending_reduction(
1341 client_order_id,
1342 bet_id.to_string(),
1343 Quantity::from(10),
1344 Quantity::from(4),
1345 );
1346
1347 assert_eq!(
1348 state.confirm_pending_reduction(&client_order_id, bet_id, active_quantity),
1349 expected,
1350 );
1351 assert_eq!(state.reduced_quantity(bet_id), expected);
1352 }
1353
1354 #[rstest]
1355 fn pending_reduction_validates_identity_and_resolves_once() {
1356 let client_order_id = ClientOrderId::from("O-1");
1357 let other_client_order_id = ClientOrderId::from("O-2");
1358 let bet_id = "bet-1";
1359 let mut state = OcmState::default();
1360 state.register_pending_reduction(
1361 client_order_id,
1362 bet_id.to_string(),
1363 Quantity::from(10),
1364 Quantity::from(4),
1365 );
1366
1367 let mismatched =
1368 state.confirm_pending_reduction(&other_client_order_id, bet_id, Quantity::from(4));
1369 assert_eq!(mismatched, None);
1370 assert!(!state.complete_pending_reduction(
1371 &other_client_order_id,
1372 bet_id,
1373 Quantity::from(4),
1374 ));
1375 state.clear_pending_reduction(&other_client_order_id, bet_id);
1376
1377 assert_eq!(
1378 state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4)),
1379 Some(Quantity::from(4)),
1380 );
1381 assert_eq!(
1382 state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4)),
1383 None,
1384 "a confirmed reduction must not resolve twice",
1385 );
1386
1387 state.clear_pending_reduction(&client_order_id, bet_id);
1388
1389 assert_eq!(state.reduced_quantity(bet_id), None);
1390 assert_eq!(
1391 state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4)),
1392 None,
1393 "a discarded reduction must not resolve from a later observation",
1394 );
1395 }
1396
1397 #[rstest]
1398 fn terminal_order_retention_evicts_pending_reduction_state() {
1399 let client_order_id = ClientOrderId::from("O-1");
1400 let bet_id = "bet-0";
1401 let mut state = OcmState::default();
1402 state.register_pending_reduction(
1403 client_order_id,
1404 bet_id.to_string(),
1405 Quantity::from(10),
1406 Quantity::from(4),
1407 );
1408 state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4));
1409
1410 for index in 0..=OcmState::DEDUP_RETENTION {
1411 state.mark_terminal_order(format!("bet-{index}"));
1412 }
1413
1414 assert_eq!(state.reduced_quantity(bet_id), None);
1415 }
1416
1417 #[rstest]
1418 fn canceled_replace_keeps_one_terminal_retention_entry() {
1419 let mut state = OcmState::default();
1420 let client_order_id = ClientOrderId::from("O-1");
1421
1422 state.mark_canceled_replace(client_order_id, "old-bet");
1423 state.retain_terminal_order(client_order_id, "old-bet");
1424
1425 assert!(state.terminal_orders.contains("old-bet"));
1426 assert_eq!(state.terminal_order_queue.len(), 1);
1427 }
1428
1429 #[rstest]
1430 fn canceled_replace_without_ocm_has_terminal_retention_entry() {
1431 let mut state = OcmState::default();
1432 let client_order_id = ClientOrderId::from("O-1");
1433
1434 state.mark_canceled_replace(client_order_id, "old-bet");
1435
1436 assert!(state.terminal_orders.contains("old-bet"));
1437 assert_eq!(state.terminal_order_queue.len(), 1);
1438 }
1439}