Skip to main content

nautilus_betfair/stream/
ocm.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Shared OCM stream handler state.
17
18use 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/// Shared mutable state for the OCM stream handler.
93///
94/// Accessed by both the TCP reader closure and the execution client methods
95/// (submit, modify, connect/disconnect). All access goes through `Arc<Mutex<>>`.
96///
97/// Terminal retention owns order correlation, per-Bet fill and reduction state, and replacement
98/// history. The first REST, OCM, or reconciliation observation with active quantity at least the
99/// requested quantity and below the original confirms a pending reduction; later observations are
100/// no-ops.
101#[derive(Clone, Debug, Default)]
102pub struct OcmState {
103    /// Tracks cumulative per-bet fill and void state for deduplication and reconciliation.
104    pub fill_tracker: FillTracker,
105    /// Bet IDs that have received a terminal event (cancel, lapse, fill-complete).
106    pub terminal_orders: AHashSet<BetId>,
107    /// Old bet IDs from replace operations, to suppress late stream updates.
108    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    /// Bounds dedup memory while retaining recent delayed stream and REST overlap.
119    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    /// Registers a customer_order_ref mapping for a new order.
132    pub fn register_customer_order_ref(&mut self, client_order_id: ClientOrderId) {
133        let _ = self.register_order_ref(client_order_id);
134    }
135
136    /// Registers both current and legacy customer_order_ref truncations.
137    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    /// Records the submitting strategy for a tracked order.
144    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    /// Returns the submitting strategy for a tracked order, if known.
237    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    /// Records that acceptance has been emitted for a tracked order.
265    ///
266    /// Returns `true` when this call newly marks the order accepted (the caller
267    /// should emit `OrderAccepted`), or `false` when acceptance was already emitted.
268    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    /// Removes customer_order_ref mappings for a client_order_id.
333    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    /// Resolves a client_order_id from the unmatched order's rfo field.
393    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    /// Returns `true` if a cancel/lapse for this bet should be suppressed
401    /// because a replace operation is pending or the bet was already replaced.
402    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    /// Returns `true` when `bet_id` is the order's most recently placed Bet.
482    pub(crate) fn is_current_venue_order_id(
483        &self,
484        client_order_id: &ClientOrderId,
485        bet_id: &str,
486    ) -> bool {
487        self.order_correlations
488            .get(client_order_id)
489            .and_then(|correlation| correlation.venue_order_id.as_ref())
490            .is_some_and(|venue_order_id| venue_order_id.as_str() == bet_id)
491    }
492
493    pub(crate) fn mark_canceled_replace(&mut self, client_order_id: ClientOrderId, bet_id: &str) {
494        self.canceled_replace_bet_ids.insert(bet_id.to_string());
495        self.retain_terminal_order(client_order_id, bet_id);
496    }
497
498    pub(crate) fn is_canceled_replace(&self, bet_id: &str) -> bool {
499        self.canceled_replace_bet_ids.contains(bet_id)
500    }
501
502    pub(crate) fn is_redundant_terminal_update(&self, order: &UnmatchedOrder) -> bool {
503        self.terminal_orders.contains(&order.id)
504            && !self.is_canceled_replace(&order.id)
505            && !self.fill_tracker.has_unseen_fill(order)
506            && !self.fill_tracker.has_unseen_fill_void(order)
507    }
508
509    pub(crate) fn clear_canceled_replace(&mut self, bet_id: &str) {
510        self.canceled_replace_bet_ids.remove(bet_id);
511    }
512
513    /// Promotes a new bet observed while a price replacement is pending.
514    ///
515    /// Returns the total order quantity and old Bet ID when `new_bet_id` differs from one pending
516    /// old Bet ID.
517    pub(crate) fn promote_pending_replace(
518        &mut self,
519        client_order_id: &ClientOrderId,
520        new_bet_id: &str,
521        replacement_quantity: Quantity,
522    ) -> Option<(Quantity, String)> {
523        let mut old_bet_ids = self
524            .pending_replace_state
525            .keys()
526            .filter(|(candidate, old_bet_id)| {
527                candidate == client_order_id && old_bet_id != new_bet_id
528            })
529            .map(|(_, old_bet_id)| old_bet_id.clone());
530
531        let old_bet_id = old_bet_ids.next()?;
532        if old_bet_ids.next().is_some() {
533            return None;
534        }
535
536        let pending = self.complete_pending_replace(
537            *client_order_id,
538            &old_bet_id,
539            VenueOrderId::from(new_bet_id),
540        )?;
541
542        Some((
543            pending.total_quantity.unwrap_or(replacement_quantity),
544            old_bet_id,
545        ))
546    }
547
548    pub(crate) fn complete_pending_replace(
549        &mut self,
550        client_order_id: ClientOrderId,
551        old_bet_id: &str,
552        new_venue_order_id: VenueOrderId,
553    ) -> Option<PendingReplaceState> {
554        if self
555            .replaced_venue_order_ids
556            .contains(new_venue_order_id.as_str())
557        {
558            return None;
559        }
560
561        let pending = self.take_pending_replace(client_order_id, old_bet_id)?;
562        self.mark_replaced_venue_order_id(client_order_id, old_bet_id.to_string());
563        self.bind_venue_order_id(&client_order_id, new_venue_order_id);
564        Some(pending)
565    }
566
567    fn mark_replaced_venue_order_id(&mut self, client_order_id: ClientOrderId, bet_id: String) {
568        self.mark_correlated_terminal_order(client_order_id, &bet_id);
569        self.replaced_venue_order_ids.insert(bet_id);
570    }
571
572    pub(crate) fn register_pending_reduction(
573        &mut self,
574        client_order_id: ClientOrderId,
575        bet_id: String,
576        original_quantity: Quantity,
577        requested_quantity: Quantity,
578    ) {
579        self.pending_reductions.insert(
580            bet_id,
581            PendingReductionState {
582                client_order_id,
583                original_quantity,
584                requested_quantity,
585                confirmed_quantity: None,
586            },
587        );
588    }
589
590    pub(crate) fn replaced_matched_quantity(&self, client_order_id: &ClientOrderId) -> Decimal {
591        self.order_correlations
592            .get(client_order_id)
593            .into_iter()
594            .flat_map(|correlation| &correlation.venue_order_ids)
595            .filter(|bet_id| self.replaced_venue_order_ids.contains(*bet_id))
596            .map(|bet_id| self.fill_tracker.matched_quantity(bet_id))
597            .sum()
598    }
599
600    pub(crate) fn confirm_pending_reduction(
601        &mut self,
602        client_order_id: &ClientOrderId,
603        bet_id: &str,
604        active_quantity: Quantity,
605    ) -> Option<Quantity> {
606        let pending = self.pending_reductions.get_mut(bet_id)?;
607
608        if !pending.can_confirm(client_order_id, active_quantity) {
609            return None;
610        }
611
612        pending.confirmed_quantity = Some(active_quantity);
613        Some(active_quantity)
614    }
615
616    pub(crate) fn complete_pending_reduction(
617        &mut self,
618        client_order_id: &ClientOrderId,
619        bet_id: &str,
620        quantity: Quantity,
621    ) -> bool {
622        let Some(pending) = self.pending_reductions.get_mut(bet_id) else {
623            return false;
624        };
625
626        if !pending.is_unconfirmed_for(client_order_id) {
627            return false;
628        }
629
630        pending.confirmed_quantity = Some(quantity);
631        true
632    }
633
634    pub(crate) fn reduced_quantity(&self, bet_id: &str) -> Option<Quantity> {
635        self.pending_reductions
636            .get(bet_id)
637            .and_then(|pending| pending.confirmed_quantity)
638    }
639
640    pub(crate) fn clear_pending_reduction(
641        &mut self,
642        client_order_id: &ClientOrderId,
643        bet_id: &str,
644    ) {
645        if self
646            .pending_reductions
647            .get(bet_id)
648            .is_some_and(|pending| pending.client_order_id == *client_order_id)
649        {
650            self.pending_reductions.remove(bet_id);
651        }
652    }
653
654    /// Retains a locally owned terminal identity and all of its Bet ID state.
655    pub(crate) fn retain_terminal_order(&mut self, client_order_id: ClientOrderId, bet_id: &str) {
656        if self.replaced_venue_order_ids.contains(bet_id)
657            || self.has_pending_replace(&client_order_id)
658        {
659            self.mark_correlated_terminal_order(client_order_id, bet_id);
660            return;
661        }
662
663        self.bind_venue_order_id(&client_order_id, VenueOrderId::from(bet_id));
664        self.mark_correlated_terminal_order(client_order_id, bet_id);
665
666        self.retain_terminal_identity(client_order_id);
667    }
668
669    fn mark_correlated_terminal_order(&mut self, client_order_id: ClientOrderId, bet_id: &str) {
670        if self
671            .pending_reductions
672            .get(bet_id)
673            .is_some_and(|pending| pending.is_unconfirmed_for(&client_order_id))
674        {
675            self.pending_reductions.remove(bet_id);
676        }
677
678        let newly_correlated = self
679            .order_correlations
680            .entry(client_order_id)
681            .or_default()
682            .venue_order_ids
683            .insert(bet_id.to_string());
684        if newly_correlated {
685            self.remove_external_terminal_retention(bet_id);
686        }
687        self.terminal_orders.insert(bet_id.to_string());
688    }
689
690    fn retain_terminal_identity(&mut self, client_order_id: ClientOrderId) {
691        let correlation = self.order_correlations.entry(client_order_id).or_default();
692        if correlation.terminal_retained {
693            return;
694        }
695
696        correlation.terminal_retained = true;
697        self.push_terminal_identity(TerminalRetentionKey::Owned(client_order_id));
698    }
699
700    /// Records an external terminal bet and bounds its stream and REST dedup state.
701    pub fn mark_terminal_order(&mut self, bet_id: String) {
702        if !self.terminal_orders.insert(bet_id.clone()) {
703            return;
704        }
705
706        self.push_terminal_identity(TerminalRetentionKey::External(bet_id));
707    }
708
709    fn push_terminal_identity(&mut self, key: TerminalRetentionKey) {
710        self.terminal_order_queue.push_back(key);
711        if self.terminal_order_queue.len() > Self::DEDUP_RETENTION
712            && let Some(expired) = self.terminal_order_queue.pop_front()
713        {
714            self.evict_terminal_identity(expired);
715        }
716    }
717
718    fn remove_terminal_retention(&mut self, client_order_id: &ClientOrderId) {
719        if self
720            .order_correlations
721            .get_mut(client_order_id)
722            .is_some_and(|correlation| std::mem::take(&mut correlation.terminal_retained))
723        {
724            self.terminal_order_queue.retain(
725                |key| !matches!(key, TerminalRetentionKey::Owned(candidate) if candidate == client_order_id),
726            );
727        }
728    }
729
730    fn remove_external_terminal_retention(&mut self, bet_id: &str) {
731        if !self.terminal_orders.contains(bet_id) {
732            return;
733        }
734
735        self.terminal_order_queue.retain(
736            |key| !matches!(key, TerminalRetentionKey::External(candidate) if candidate == bet_id),
737        );
738    }
739
740    pub(crate) fn mark_order_active(&mut self, client_order_id: &ClientOrderId, bet_id: &str) {
741        self.remove_terminal_retention(client_order_id);
742        self.terminal_orders.remove(bet_id);
743        self.canceled_replace_bet_ids.remove(bet_id);
744    }
745
746    fn evict_terminal_identity(&mut self, key: TerminalRetentionKey) {
747        match key {
748            TerminalRetentionKey::Owned(client_order_id) => {
749                let retained = self
750                    .order_correlations
751                    .get_mut(&client_order_id)
752                    .is_some_and(|correlation| std::mem::take(&mut correlation.terminal_retained));
753                if retained {
754                    self.remove_order_correlation(&client_order_id);
755                }
756            }
757            TerminalRetentionKey::External(bet_id) => {
758                self.remove_venue_order_state(&bet_id);
759            }
760        }
761    }
762
763    fn remove_venue_order_state(&mut self, bet_id: &str) {
764        self.terminal_orders.remove(bet_id);
765        self.replaced_venue_order_ids.remove(bet_id);
766        self.canceled_replace_bet_ids.remove(bet_id);
767        self.pending_reductions.remove(bet_id);
768        self.fill_tracker.prune(bet_id);
769    }
770
771    /// Anchors the fill tracker against cached orders so the post-reconnect
772    /// image neither treats cumulative size as a new fill nor re-emits a
773    /// fill that was published via another channel.
774    pub fn sync_from_orders(&mut self, orders: &[OrderSyncEntry]) {
775        for entry in orders {
776            self.restore_order(
777                entry.client_order_id,
778                entry.strategy_id,
779                VenueOrderId::from(entry.bet_id.as_str()),
780            );
781
782            for venue_order_id in &entry.venue_order_ids {
783                if venue_order_id != &entry.bet_id {
784                    self.mark_replaced_venue_order_id(
785                        entry.client_order_id,
786                        venue_order_id.clone(),
787                    );
788                }
789            }
790
791            if entry.is_closed {
792                self.retain_terminal_order(entry.client_order_id, &entry.bet_id);
793            } else {
794                self.mark_order_active(&entry.client_order_id, &entry.bet_id);
795            }
796
797            if entry.filled_qty > Decimal::ZERO {
798                self.fill_tracker
799                    .sync_order(&entry.bet_id, entry.filled_qty, entry.avg_px);
800            }
801
802            if !entry.trade_ids.is_empty() {
803                self.fill_tracker
804                    .seed_published_trade_ids(entry.trade_ids.iter().cloned());
805            }
806        }
807    }
808}
809
810#[cfg(test)]
811mod tests {
812    use rstest::rstest;
813
814    use super::*;
815
816    #[rstest]
817    fn owned_terminal_retention_evicts_all_identity_state_together() {
818        let mut state = OcmState::default();
819        let client_order_id = ClientOrderId::from("O-TERMINAL-0");
820        let strategy_id = StrategyId::from("S-001");
821        let current_bet_id = "current-bet-0";
822        let replaced_bet_id = "replaced-bet-0";
823        let size_matched = Decimal::from(2);
824        let average_price = Decimal::from(3);
825        state.restore_order(
826            client_order_id,
827            strategy_id,
828            VenueOrderId::from(current_bet_id),
829        );
830        state.mark_replaced_venue_order_id(client_order_id, replaced_bet_id.to_string());
831        assert!(state.is_accepted(&client_order_id));
832        state.register_pending_reduction(
833            client_order_id,
834            current_bet_id.to_string(),
835            Quantity::from(10),
836            Quantity::from(4),
837        );
838        state.confirm_pending_reduction(&client_order_id, current_bet_id, Quantity::from(4));
839        state.fill_tracker.advance_cumulative_fill(
840            current_bet_id,
841            size_matched,
842            Some(average_price),
843            average_price,
844        );
845        state.fill_tracker.advance_cumulative_fill(
846            replaced_bet_id,
847            size_matched,
848            Some(average_price),
849            average_price,
850        );
851        state.retain_terminal_order(client_order_id, current_bet_id);
852
853        for index in 0..OcmState::DEDUP_RETENTION {
854            state.mark_terminal_order(format!("external-bet-{index}"));
855        }
856
857        let current_replay = state.fill_tracker.advance_cumulative_fill(
858            current_bet_id,
859            size_matched,
860            Some(average_price),
861            average_price,
862        );
863        let replaced_replay = state.fill_tracker.advance_cumulative_fill(
864            replaced_bet_id,
865            size_matched,
866            Some(average_price),
867            average_price,
868        );
869        let customer_order_ref = make_customer_order_ref(client_order_id.as_str());
870
871        assert_eq!(state.order_strategy_id(&client_order_id), None);
872        assert_eq!(
873            state.resolve_client_order_id(Some(&customer_order_ref)),
874            None,
875        );
876        assert_eq!(state.reduced_quantity(current_bet_id), None);
877        assert!(!state.terminal_orders.contains(current_bet_id));
878        assert!(!state.terminal_orders.contains(replaced_bet_id));
879        assert!(!state.replaced_venue_order_ids.contains(replaced_bet_id));
880        assert!(current_replay.is_some());
881        assert!(replaced_replay.is_some());
882    }
883
884    #[rstest]
885    fn terminal_retention_does_not_evict_active_or_ambiguous_replace_identity() {
886        let mut state = OcmState::default();
887        let active_client_order_id = ClientOrderId::from("O-ACTIVE");
888        let ambiguous_client_order_id = ClientOrderId::from("O-AMBIGUOUS");
889        let strategy_id = StrategyId::from("S-001");
890        state.restore_order(
891            active_client_order_id,
892            strategy_id,
893            VenueOrderId::from("active-bet"),
894        );
895        state.restore_order(
896            ambiguous_client_order_id,
897            strategy_id,
898            VenueOrderId::from("ambiguous-bet"),
899        );
900        state.register_pending_replace(
901            ambiguous_client_order_id,
902            "ambiguous-bet".to_string(),
903            Some(Quantity::from(10)),
904        );
905        state.mark_pending_replace_ambiguous(ambiguous_client_order_id, "ambiguous-bet");
906        state.retain_terminal_order(ambiguous_client_order_id, "ambiguous-bet");
907
908        for index in 0..=OcmState::DEDUP_RETENTION {
909            state.mark_terminal_order(format!("external-bet-{index}"));
910        }
911
912        assert_eq!(
913            state.order_strategy_id(&active_client_order_id),
914            Some(strategy_id),
915        );
916        assert!(!state.mark_accepted(active_client_order_id));
917        assert_eq!(
918            state.order_strategy_id(&ambiguous_client_order_id),
919            Some(strategy_id),
920        );
921        assert!(state.pending_replace_awaits_reconciliation(
922            &ambiguous_client_order_id,
923            "ambiguous-bet",
924        ));
925        assert!(state.should_suppress_cancel(&ambiguous_client_order_id, "ambiguous-bet",));
926    }
927
928    #[rstest]
929    fn successful_replace_without_old_leg_ocm_uses_terminal_identity_lifecycle() {
930        let mut state = OcmState::default();
931        let client_order_id = ClientOrderId::from("O-REPLACE");
932        let strategy_id = StrategyId::from("S-001");
933        state.restore_order(client_order_id, strategy_id, VenueOrderId::from("old-bet"));
934        state.register_pending_replace(
935            client_order_id,
936            "old-bet".to_string(),
937            Some(Quantity::from(10)),
938        );
939
940        let completed = state.complete_pending_replace(
941            client_order_id,
942            "old-bet",
943            VenueOrderId::from("new-bet"),
944        );
945
946        assert!(completed.is_some());
947        assert!(state.pending_replace_state.is_empty());
948        assert!(state.replaced_venue_order_ids.contains("old-bet"));
949        assert!(state.terminal_orders.contains("old-bet"));
950        assert_eq!(
951            state.client_order_id_by_venue_order_id("old-bet"),
952            Some(client_order_id),
953        );
954        assert_eq!(
955            state.client_order_id_by_venue_order_id("new-bet"),
956            Some(client_order_id),
957        );
958        assert!(!state.is_retained_terminal_order(&client_order_id));
959
960        state.retain_terminal_order(client_order_id, "old-bet");
961
962        assert_eq!(
963            state
964                .order_correlations
965                .get(&client_order_id)
966                .and_then(|correlation| correlation.venue_order_id),
967            Some(VenueOrderId::from("new-bet")),
968        );
969        assert!(!state.is_retained_terminal_order(&client_order_id));
970
971        state.retain_terminal_order(client_order_id, "new-bet");
972        state.retain_terminal_order(client_order_id, "old-bet");
973
974        assert_eq!(
975            state
976                .order_correlations
977                .get(&client_order_id)
978                .and_then(|correlation| correlation.venue_order_id),
979            Some(VenueOrderId::from("new-bet")),
980        );
981        assert_eq!(state.terminal_order_queue.len(), 1);
982    }
983
984    #[rstest]
985    fn replace_resolution_does_not_discard_another_pending_mutation() {
986        let mut state = OcmState::default();
987        let client_order_id = ClientOrderId::from("O-REPLACE");
988        state.restore_order(
989            client_order_id,
990            StrategyId::from("S-001"),
991            VenueOrderId::from("old-bet-1"),
992        );
993        state.register_pending_replace(
994            client_order_id,
995            "old-bet-1".to_string(),
996            Some(Quantity::from(10)),
997        );
998        state.register_pending_replace(
999            client_order_id,
1000            "old-bet-2".to_string(),
1001            Some(Quantity::from(10)),
1002        );
1003
1004        let unresolved_promotion =
1005            state.promote_pending_replace(&client_order_id, "new-bet", Quantity::from(10));
1006        let completed = state.complete_pending_replace(
1007            client_order_id,
1008            "old-bet-1",
1009            VenueOrderId::from("new-bet"),
1010        );
1011
1012        assert_eq!(unresolved_promotion, None);
1013        assert!(completed.is_some());
1014        assert!(
1015            !state
1016                .pending_replace_state
1017                .contains_key(&(client_order_id, "old-bet-1".to_string()))
1018        );
1019        assert!(
1020            state
1021                .pending_replace_state
1022                .contains_key(&(client_order_id, "old-bet-2".to_string()))
1023        );
1024        assert!(state.replaced_venue_order_ids.contains("old-bet-1"));
1025        assert!(!state.replaced_venue_order_ids.contains("old-bet-2"));
1026        assert_eq!(
1027            state.client_order_id_by_venue_order_id("new-bet"),
1028            Some(client_order_id),
1029        );
1030    }
1031
1032    #[rstest]
1033    fn historical_bet_migrates_from_external_to_owned_identity() {
1034        let mut state = OcmState::default();
1035        let client_order_id = ClientOrderId::from("O-MIGRATED");
1036        let old_bet_id = "old-bet";
1037        let size_matched = Decimal::from(2);
1038        let average_price = Decimal::from(3);
1039        state.mark_terminal_order(old_bet_id.to_string());
1040        state.fill_tracker.advance_cumulative_fill(
1041            old_bet_id,
1042            size_matched,
1043            Some(average_price),
1044            average_price,
1045        );
1046        state.restore_order(
1047            client_order_id,
1048            StrategyId::from("S-001"),
1049            VenueOrderId::from("current-bet"),
1050        );
1051        state.mark_replaced_venue_order_id(client_order_id, old_bet_id.to_string());
1052
1053        for index in 0..=OcmState::DEDUP_RETENTION {
1054            state.mark_terminal_order(format!("external-bet-{index}"));
1055        }
1056
1057        let replay = state.fill_tracker.advance_cumulative_fill(
1058            old_bet_id,
1059            size_matched,
1060            Some(average_price),
1061            average_price,
1062        );
1063
1064        assert_eq!(
1065            state.client_order_id_by_venue_order_id(old_bet_id),
1066            Some(client_order_id),
1067        );
1068        assert!(state.terminal_orders.contains(old_bet_id));
1069        assert!(state.replaced_venue_order_ids.contains(old_bet_id));
1070        assert!(replay.is_none());
1071    }
1072
1073    #[rstest]
1074    fn sustained_owned_terminal_history_stays_bounded() {
1075        let mut state = OcmState::default();
1076        let strategy_id = StrategyId::from("S-001");
1077
1078        for index in 0..=OcmState::DEDUP_RETENTION {
1079            let client_order_id = ClientOrderId::from(format!("O-{index}"));
1080            let current_bet_id = format!("current-bet-{index}");
1081            state.restore_order(
1082                client_order_id,
1083                strategy_id,
1084                VenueOrderId::from(current_bet_id.as_str()),
1085            );
1086            state.mark_replaced_venue_order_id(client_order_id, format!("replaced-bet-{index}"));
1087            state.retain_terminal_order(client_order_id, &current_bet_id);
1088        }
1089
1090        assert_eq!(state.order_correlations.len(), OcmState::DEDUP_RETENTION);
1091        assert_eq!(
1092            state
1093                .order_correlations
1094                .values()
1095                .filter(|correlation| correlation.terminal_retained)
1096                .count(),
1097            OcmState::DEDUP_RETENTION,
1098        );
1099        assert_eq!(state.terminal_order_queue.len(), OcmState::DEDUP_RETENTION);
1100        assert_eq!(
1101            state.replaced_venue_order_ids.len(),
1102            OcmState::DEDUP_RETENTION,
1103        );
1104        assert_eq!(state.terminal_orders.len(), OcmState::DEDUP_RETENTION * 2);
1105        assert_eq!(state.order_strategy_id(&ClientOrderId::from("O-0")), None,);
1106        assert_eq!(
1107            state.order_strategy_id(&ClientOrderId::from(format!(
1108                "O-{}",
1109                OcmState::DEDUP_RETENTION
1110            ))),
1111            Some(strategy_id),
1112        );
1113    }
1114
1115    #[rstest]
1116    fn terminal_order_retention_evicts_fill_tracker_state() {
1117        let mut state = OcmState::default();
1118        let first_bet_id = "bet-0";
1119        let size_matched = Decimal::new(10, 0);
1120        let average_price = Decimal::new(20, 1);
1121
1122        assert!(
1123            state
1124                .fill_tracker
1125                .advance_cumulative_fill(
1126                    first_bet_id,
1127                    size_matched,
1128                    Some(average_price),
1129                    average_price,
1130                )
1131                .is_some(),
1132        );
1133        state
1134            .replaced_venue_order_ids
1135            .insert(first_bet_id.to_string());
1136        state.mark_terminal_order(first_bet_id.to_string());
1137        for index in 1..=OcmState::DEDUP_RETENTION {
1138            state.mark_terminal_order(format!("bet-{index}"));
1139        }
1140
1141        let replay_after_eviction = state.fill_tracker.advance_cumulative_fill(
1142            first_bet_id,
1143            size_matched,
1144            Some(average_price),
1145            average_price,
1146        );
1147
1148        assert!(!state.terminal_orders.contains(first_bet_id));
1149        assert!(!state.replaced_venue_order_ids.contains(first_bet_id));
1150        assert!(replay_after_eviction.is_some());
1151    }
1152
1153    #[rstest]
1154    fn active_submission_collision_is_rejected() {
1155        let suffix = "12345678901234567890123456789012";
1156        let first = ClientOrderId::from(format!("FIRST-{suffix}"));
1157        let second = ClientOrderId::from(format!("SECOND-{suffix}"));
1158        let strategy_id = StrategyId::from("S-001");
1159        let mut state = OcmState::default();
1160
1161        state.register_submission(first, strategy_id).unwrap();
1162        let collision = state.register_submission(second, strategy_id);
1163
1164        assert_eq!(collision, Err(suffix.to_string()));
1165        assert_eq!(state.resolve_client_order_id(Some(suffix)), Some(first));
1166        assert_eq!(state.order_strategy_id(&second), None);
1167    }
1168
1169    #[rstest]
1170    fn submission_collision_with_restored_legacy_reference_is_rejected() {
1171        let reference = "12345678901234567890123456789012";
1172        let restored = ClientOrderId::from(format!("{reference}-RESTORED"));
1173        let fresh = ClientOrderId::from(format!("FRESH-{reference}"));
1174        let strategy_id = StrategyId::from("S-001");
1175        let mut state = OcmState::default();
1176
1177        state.restore_order(restored, strategy_id, VenueOrderId::from("bet-restored"));
1178        let collision = state.register_submission(fresh, strategy_id);
1179
1180        assert_eq!(collision, Err(reference.to_string()));
1181        assert_eq!(
1182            state.resolve_client_order_id(Some(reference)),
1183            Some(restored),
1184        );
1185        assert_eq!(state.order_strategy_id(&fresh), None);
1186    }
1187
1188    #[rstest]
1189    fn restored_current_reference_ambiguity_recovers_after_cleanup() {
1190        let suffix = "12345678901234567890123456789012";
1191        let first = ClientOrderId::from(format!("FIRST-{suffix}"));
1192        let second = ClientOrderId::from(format!("SECOND-{suffix}"));
1193        let strategy_id = StrategyId::from("S-001");
1194        let mut state = OcmState::default();
1195
1196        state.restore_order(first, strategy_id, VenueOrderId::from("bet-1"));
1197        state.restore_order(second, strategy_id, VenueOrderId::from("bet-2"));
1198
1199        assert_eq!(
1200            state.customer_order_ref_resolution(suffix),
1201            Some(CustomerOrderRefResolution::Ambiguous),
1202        );
1203        assert_eq!(state.resolve_client_order_id(Some(suffix)), None);
1204        state.remove_order_correlation(&first);
1205        assert_eq!(state.resolve_client_order_id(Some(suffix)), Some(second));
1206        assert_eq!(
1207            state
1208                .order_correlations
1209                .get(&second)
1210                .and_then(|correlation| correlation.venue_order_id),
1211            Some(VenueOrderId::from("bet-2")),
1212        );
1213        assert!(!state.mark_accepted(second));
1214    }
1215
1216    #[rstest]
1217    fn restored_legacy_reference_ambiguity_recovers_after_cleanup() {
1218        let prefix = "12345678901234567890123456789012";
1219        let first = ClientOrderId::from(format!("{prefix}-FIRST"));
1220        let second = ClientOrderId::from(format!("{prefix}-SECOND"));
1221        let strategy_id = StrategyId::from("S-001");
1222        let mut state = OcmState::default();
1223
1224        state.restore_order(first, strategy_id, VenueOrderId::from("bet-1"));
1225        state.restore_order(second, strategy_id, VenueOrderId::from("bet-2"));
1226
1227        assert_eq!(
1228            state.customer_order_ref_resolution(prefix),
1229            Some(CustomerOrderRefResolution::Ambiguous),
1230        );
1231        assert_eq!(state.resolve_client_order_id(Some(prefix)), None);
1232        state.remove_order_correlation(&second);
1233        assert_eq!(state.resolve_client_order_id(Some(prefix)), Some(first));
1234    }
1235
1236    #[rstest]
1237    fn venue_identity_can_be_bound_and_replaced() {
1238        let client_order_id = ClientOrderId::from("O-1");
1239        let mut state = OcmState::default();
1240        state
1241            .register_submission(client_order_id, StrategyId::from("S-001"))
1242            .unwrap();
1243
1244        state.bind_venue_order_id(&client_order_id, VenueOrderId::from("bet-1"));
1245        assert_eq!(
1246            state
1247                .order_correlations
1248                .get(&client_order_id)
1249                .and_then(|correlation| correlation.venue_order_id),
1250            Some(VenueOrderId::from("bet-1")),
1251        );
1252
1253        state.bind_venue_order_id(&client_order_id, VenueOrderId::from("bet-2"));
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-2")),
1260        );
1261    }
1262
1263    #[rstest]
1264    fn claimed_acceptance_does_not_replace_restored_venue_identity() {
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("bet-current"),
1271        );
1272
1273        let claimed = state.claim_acceptance(client_order_id, VenueOrderId::from("bet-stale"));
1274
1275        assert!(!claimed);
1276        assert_eq!(
1277            state
1278                .order_correlations
1279                .get(&client_order_id)
1280                .and_then(|correlation| correlation.venue_order_id),
1281            Some(VenueOrderId::from("bet-current")),
1282        );
1283    }
1284
1285    #[rstest]
1286    fn pending_replace_promotes_only_a_different_bet() {
1287        let client_order_id = ClientOrderId::from("O-1");
1288        let mut state = OcmState::default();
1289        state.restore_order(
1290            client_order_id,
1291            StrategyId::from("S-001"),
1292            VenueOrderId::from("old-bet"),
1293        );
1294        state.register_pending_replace(
1295            client_order_id,
1296            "old-bet".to_string(),
1297            Some(Quantity::from(10)),
1298        );
1299
1300        assert_eq!(
1301            state.promote_pending_replace(&client_order_id, "old-bet", Quantity::from(8)),
1302            None
1303        );
1304        assert_eq!(
1305            state.promote_pending_replace(&client_order_id, "new-bet", Quantity::from(8)),
1306            Some((Quantity::from(10), "old-bet".to_string())),
1307        );
1308        assert!(state.pending_replace_state.is_empty());
1309        assert!(state.replaced_venue_order_ids.contains("old-bet"));
1310        assert_eq!(
1311            state
1312                .order_correlations
1313                .get(&client_order_id)
1314                .and_then(|correlation| correlation.venue_order_id),
1315            Some(VenueOrderId::from("new-bet")),
1316        );
1317        assert_eq!(
1318            state.promote_pending_replace(&client_order_id, "newer-bet", Quantity::from(8)),
1319            None,
1320        );
1321    }
1322
1323    #[rstest]
1324    fn pending_replace_does_not_promote_a_historical_bet() {
1325        let client_order_id = ClientOrderId::from("O-1");
1326        let mut state = OcmState::default();
1327        state
1328            .replaced_venue_order_ids
1329            .insert("historical-bet".to_string());
1330        state.register_pending_replace(
1331            client_order_id,
1332            "current-bet".to_string(),
1333            Some(Quantity::from(10)),
1334        );
1335
1336        assert_eq!(
1337            state.promote_pending_replace(&client_order_id, "historical-bet", Quantity::from(8)),
1338            None,
1339        );
1340        assert!(
1341            state.should_suppress_cancel(&client_order_id, "current-bet"),
1342            "historical traffic must not consume the pending replace",
1343        );
1344        assert_eq!(
1345            state.promote_pending_replace(&client_order_id, "replacement-bet", Quantity::from(8)),
1346            Some((Quantity::from(10), "current-bet".to_string())),
1347        );
1348    }
1349
1350    #[rstest]
1351    #[case::below_requested(Quantity::from(3), None)]
1352    #[case::at_requested(Quantity::from(4), Some(Quantity::from(4)))]
1353    #[case::inside_window(Quantity::from(9), Some(Quantity::from(9)))]
1354    #[case::at_original(Quantity::from(10), None)]
1355    fn pending_reduction_confirms_only_a_definitive_reduction(
1356        #[case] active_quantity: Quantity,
1357        #[case] expected: Option<Quantity>,
1358    ) {
1359        let client_order_id = ClientOrderId::from("O-1");
1360        let bet_id = "bet-1";
1361        let mut state = OcmState::default();
1362        state.register_pending_reduction(
1363            client_order_id,
1364            bet_id.to_string(),
1365            Quantity::from(10),
1366            Quantity::from(4),
1367        );
1368
1369        assert_eq!(
1370            state.confirm_pending_reduction(&client_order_id, bet_id, active_quantity),
1371            expected,
1372        );
1373        assert_eq!(state.reduced_quantity(bet_id), expected);
1374    }
1375
1376    #[rstest]
1377    fn pending_reduction_validates_identity_and_resolves_once() {
1378        let client_order_id = ClientOrderId::from("O-1");
1379        let other_client_order_id = ClientOrderId::from("O-2");
1380        let bet_id = "bet-1";
1381        let mut state = OcmState::default();
1382        state.register_pending_reduction(
1383            client_order_id,
1384            bet_id.to_string(),
1385            Quantity::from(10),
1386            Quantity::from(4),
1387        );
1388
1389        let mismatched =
1390            state.confirm_pending_reduction(&other_client_order_id, bet_id, Quantity::from(4));
1391        assert_eq!(mismatched, None);
1392        assert!(!state.complete_pending_reduction(
1393            &other_client_order_id,
1394            bet_id,
1395            Quantity::from(4),
1396        ));
1397        state.clear_pending_reduction(&other_client_order_id, bet_id);
1398
1399        assert_eq!(
1400            state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4)),
1401            Some(Quantity::from(4)),
1402        );
1403        assert_eq!(
1404            state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4)),
1405            None,
1406            "a confirmed reduction must not resolve twice",
1407        );
1408
1409        state.clear_pending_reduction(&client_order_id, bet_id);
1410
1411        assert_eq!(state.reduced_quantity(bet_id), None);
1412        assert_eq!(
1413            state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4)),
1414            None,
1415            "a discarded reduction must not resolve from a later observation",
1416        );
1417    }
1418
1419    #[rstest]
1420    fn terminal_order_retention_evicts_pending_reduction_state() {
1421        let client_order_id = ClientOrderId::from("O-1");
1422        let bet_id = "bet-0";
1423        let mut state = OcmState::default();
1424        state.register_pending_reduction(
1425            client_order_id,
1426            bet_id.to_string(),
1427            Quantity::from(10),
1428            Quantity::from(4),
1429        );
1430        state.confirm_pending_reduction(&client_order_id, bet_id, Quantity::from(4));
1431
1432        for index in 0..=OcmState::DEDUP_RETENTION {
1433            state.mark_terminal_order(format!("bet-{index}"));
1434        }
1435
1436        assert_eq!(state.reduced_quantity(bet_id), None);
1437    }
1438
1439    #[rstest]
1440    fn canceled_replace_keeps_one_terminal_retention_entry() {
1441        let mut state = OcmState::default();
1442        let client_order_id = ClientOrderId::from("O-1");
1443
1444        state.mark_canceled_replace(client_order_id, "old-bet");
1445        state.retain_terminal_order(client_order_id, "old-bet");
1446
1447        assert!(state.terminal_orders.contains("old-bet"));
1448        assert_eq!(state.terminal_order_queue.len(), 1);
1449    }
1450
1451    #[rstest]
1452    fn canceled_replace_without_ocm_has_terminal_retention_entry() {
1453        let mut state = OcmState::default();
1454        let client_order_id = ClientOrderId::from("O-1");
1455
1456        state.mark_canceled_replace(client_order_id, "old-bet");
1457
1458        assert!(state.terminal_orders.contains("old-bet"));
1459        assert_eq!(state.terminal_order_queue.len(), 1);
1460    }
1461}