Skip to main content

nautilus_okx/websocket/
dispatch.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//! WebSocket message dispatch for the OKX execution client.
17//!
18//! Routes incoming [`OKXWsMessage`] variants to the appropriate parsing and
19//! event emission paths. Tracked orders (submitted through this client) produce
20//! proper order events; untracked orders fall back to execution reports for
21//! downstream reconciliation.
22
23use std::{collections::VecDeque, fmt::Debug, hash::Hash, sync::Arc};
24
25use ahash::AHashMap;
26use dashmap::DashMap;
27use nautilus_common::cache::fifo::{FifoCache, FifoCacheMap};
28use nautilus_core::{AtomicMap, UUID4, UnixNanos, time::AtomicTime};
29use nautilus_live::{
30    ExecutionEventEmitter,
31    execution::{
32        context::{OrderContext, OrderIdentity},
33        failure::CommandFailure,
34    },
35};
36use nautilus_model::{
37    enums::{AccountType, OrderStatus},
38    events::{
39        OrderAccepted, OrderCanceled, OrderEventAny, OrderFilled, OrderRejected, OrderTriggered,
40        OrderUpdated,
41    },
42    identifiers::{
43        AccountId, ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, VenueOrderId,
44    },
45    instruments::{Instrument, InstrumentAny},
46    orders::TRIGGERABLE_ORDER_TYPES,
47    reports::FillReport,
48    types::Currency,
49};
50use parking_lot::Mutex;
51use ustr::Ustr;
52
53use crate::{
54    common::{
55        consts::{
56            OKX_FIELD_CLORDID, OKX_FIELD_SCODE, OKX_FIELD_SMSG, OKX_FIELD_SUBCODE,
57            OKX_POST_ONLY_CANCEL_REASON, OKX_POST_ONLY_CANCEL_SOURCE, OKX_SUCCESS_CODE,
58        },
59        enums::{OKXAlgoOrderStatus, OKXAlgoOrderType, OKXOrderStatus, OKXOrderType},
60        failure::{classify_okx_venue_code, classify_okx_ws_failure},
61        parse::{
62            is_market_price, parse_client_order_id, parse_millisecond_timestamp, parse_price,
63            parse_quantity,
64        },
65    },
66    http::models::{OKXAccount, OKXCancelAlgoOrderResponse, OKXPosition, OKXSpreadOrder},
67    websocket::{
68        client::PendingOrderInfo,
69        enums::OKXWsOperation,
70        handler::{is_post_only_auto_cancel, is_unfilled_rpi_cancel},
71        messages::{ExecutionReport, OKXAlgoOrderMsg, OKXOrderMsg, OKXWsMessage},
72        parse::{
73            FeeCache, FilledQtyCache, OrderStateSnapshot, ParsedOrderEvent,
74            is_terminal_order_state, parse_algo_order_msg, parse_algo_order_status_report,
75            parse_order_event, parse_order_msg, parse_spread_order_event, parse_spread_order_msg,
76            update_fee_fill_caches,
77        },
78    },
79};
80
81/// Maximum entries held by the dedup sets before the oldest is evicted.
82const DEDUP_CAPACITY: usize = 10_000;
83
84#[derive(Clone, Copy, Debug, PartialEq, Eq)]
85struct OrderVenueBinding {
86    parent: VenueOrderId,
87    child: Option<VenueOrderId>,
88}
89
90#[derive(Debug)]
91struct OrderLifecycleBindings {
92    client_by_parent: AHashMap<VenueOrderId, ClientOrderId>,
93    venue_by_client: AHashMap<ClientOrderId, OrderVenueBinding>,
94    terminal_client_by_parent: FifoCacheMap<VenueOrderId, ClientOrderId, DEDUP_CAPACITY>,
95}
96
97impl Default for OrderLifecycleBindings {
98    fn default() -> Self {
99        Self {
100            client_by_parent: AHashMap::new(),
101            venue_by_client: AHashMap::new(),
102            terminal_client_by_parent: FifoCacheMap::new(),
103        }
104    }
105}
106
107impl OrderLifecycleBindings {
108    fn client_order_id(&self, parent: &VenueOrderId) -> Option<ClientOrderId> {
109        self.client_by_parent
110            .get(parent)
111            .or_else(|| self.terminal_client_by_parent.get(parent))
112            .copied()
113    }
114
115    fn finish(&mut self, client_order_id: ClientOrderId, binding: OrderVenueBinding) {
116        self.client_by_parent.remove(&binding.parent);
117        self.venue_by_client.remove(&client_order_id);
118        self.terminal_client_by_parent
119            .insert(binding.parent, client_order_id);
120    }
121}
122
123#[derive(Clone, Copy, Debug, PartialEq, Eq)]
124enum ExecutionUpdateRoute {
125    Tracked(ClientOrderId, OrderContext),
126    External,
127    Suppressed,
128}
129
130#[derive(Clone, Copy, Debug, PartialEq, Eq)]
131enum LinkedChildResolution {
132    Bound(ClientOrderId),
133    Held,
134    External,
135}
136
137#[derive(Debug)]
138struct PendingLinkedChild {
139    parent_venue_order_id: VenueOrderId,
140    candidate_client_order_ids: Vec<ClientOrderId>,
141    message: OKXOrderMsg,
142}
143
144#[derive(Debug)]
145struct DedupCache<K>
146where
147    K: Clone + Debug + Eq + Hash,
148{
149    inner: Mutex<FifoCache<K, DEDUP_CAPACITY>>,
150}
151
152impl<K> DedupCache<K>
153where
154    K: Clone + Debug + Eq + Hash,
155{
156    fn new() -> Self {
157        Self {
158            inner: Mutex::new(FifoCache::new()),
159        }
160    }
161
162    fn contains(&self, key: &K) -> bool {
163        self.inner.lock().contains(key)
164    }
165
166    fn insert(&self, key: K) -> bool {
167        self.inner.lock().insert(key)
168    }
169
170    fn remove(&self, key: &K) {
171        self.inner.lock().remove(key);
172    }
173}
174
175/// Shared state for cross-stream event deduplication between the private
176/// and business WebSocket dispatch loops.
177#[derive(Debug)]
178pub struct WsDispatchState {
179    pub order_identities: DashMap<ClientOrderId, OrderIdentity>,
180    order_contexts: DashMap<ClientOrderId, OrderContext>,
181    pub(crate) pending_orders: Arc<DashMap<String, PendingOrderInfo>>,
182    pub(crate) pending_cancels: Arc<DashMap<String, PendingOrderInfo>>,
183    pub(crate) pending_amends: Arc<DashMap<String, PendingOrderInfo>>,
184    accepted_venue_order_ids: Mutex<FifoCacheMap<ClientOrderId, VenueOrderId, DEDUP_CAPACITY>>,
185    triggered_orders: DedupCache<ClientOrderId>,
186    filled_orders: DedupCache<ClientOrderId>,
187    terminal_orders: DedupCache<ClientOrderId>,
188    emitted_trades: DedupCache<TradeId>,
189    post_only_rejections: DedupCache<Ustr>,
190    lifecycle_bindings: Mutex<OrderLifecycleBindings>,
191    pending_linked_children: Mutex<VecDeque<PendingLinkedChild>>,
192    linked_child_notify: tokio::sync::Notify,
193}
194
195impl Default for WsDispatchState {
196    fn default() -> Self {
197        Self {
198            order_identities: DashMap::new(),
199            order_contexts: DashMap::new(),
200            pending_orders: Arc::new(DashMap::new()),
201            pending_cancels: Arc::new(DashMap::new()),
202            pending_amends: Arc::new(DashMap::new()),
203            accepted_venue_order_ids: Mutex::new(FifoCacheMap::new()),
204            triggered_orders: DedupCache::new(),
205            filled_orders: DedupCache::new(),
206            terminal_orders: DedupCache::new(),
207            emitted_trades: DedupCache::new(),
208            post_only_rejections: DedupCache::new(),
209            lifecycle_bindings: Mutex::new(OrderLifecycleBindings::default()),
210            pending_linked_children: Mutex::new(VecDeque::new()),
211            linked_child_notify: tokio::sync::Notify::new(),
212        }
213    }
214}
215
216impl WsDispatchState {
217    // Creates a dispatch state sharing the pending operation maps
218    // with the WebSocket client that populates them
219    pub(crate) fn with_pending_maps(
220        pending_orders: Arc<DashMap<String, PendingOrderInfo>>,
221        pending_cancels: Arc<DashMap<String, PendingOrderInfo>>,
222        pending_amends: Arc<DashMap<String, PendingOrderInfo>>,
223    ) -> Self {
224        Self {
225            pending_orders,
226            pending_cancels,
227            pending_amends,
228            ..Default::default()
229        }
230    }
231
232    pub(crate) fn track_order_context(&self, context: OrderContext) {
233        self.order_contexts
234            .insert(context.identity.client_order_id, context);
235    }
236
237    pub(crate) fn bind_algo_parent(
238        &self,
239        client_order_id: ClientOrderId,
240        venue_order_id: VenueOrderId,
241    ) {
242        let mut bindings = self.lifecycle_bindings.lock();
243        let binding = bindings.venue_by_client.get(&client_order_id).copied();
244        if binding.is_some_and(|binding| binding.parent != venue_order_id) {
245            log::error!(
246                "Ignoring conflicting algo parent binding for {client_order_id}: expected={:?} received={venue_order_id}",
247                binding.map(|binding| binding.parent),
248            );
249            return;
250        }
251
252        bindings
253            .client_by_parent
254            .insert(venue_order_id, client_order_id);
255        bindings.venue_by_client.insert(
256            client_order_id,
257            binding.unwrap_or(OrderVenueBinding {
258                parent: venue_order_id,
259                child: None,
260            }),
261        );
262        self.linked_child_notify.notify_one();
263    }
264
265    pub(crate) fn order_venue_binding(
266        &self,
267        client_order_id: ClientOrderId,
268    ) -> Option<(VenueOrderId, bool)> {
269        let bindings = self.lifecycle_bindings.lock();
270        bindings
271            .venue_by_client
272            .get(&client_order_id)
273            .map(|binding| {
274                (
275                    binding.child.unwrap_or(binding.parent),
276                    binding.child.is_some(),
277                )
278            })
279    }
280
281    pub(crate) fn order_identity(&self, client_order_id: ClientOrderId) -> Option<OrderIdentity> {
282        self.order_contexts
283            .get(&client_order_id)
284            .map(|entry| entry.identity)
285            .or_else(|| {
286                self.order_identities
287                    .get(&client_order_id)
288                    .map(|entry| *entry)
289            })
290    }
291
292    pub(crate) fn remove_order_tracking(&self, client_order_id: ClientOrderId) {
293        let pending = self.pending_linked_children.lock();
294        let removed_context = self.order_contexts.remove(&client_order_id).is_some();
295        self.order_identities.remove(&client_order_id);
296        drop(pending);
297
298        if removed_context {
299            self.linked_child_notify.notify_one();
300        }
301    }
302
303    pub(crate) fn resolve_algo_submit_failure(
304        &self,
305        client_order_id: ClientOrderId,
306        failure: &CommandFailure,
307    ) {
308        let bindings = self.lifecycle_bindings.lock();
309        if matches!(failure, CommandFailure::Ambiguous(_))
310            && bindings.venue_by_client.contains_key(&client_order_id)
311        {
312            return;
313        }
314
315        self.remove_order_tracking(client_order_id);
316    }
317
318    pub(crate) async fn wait_for_linked_child_route(&self) {
319        self.linked_child_notify.notified().await;
320    }
321
322    fn resolve_or_hold_linked_child(
323        &self,
324        parent_venue_order_id: VenueOrderId,
325        instrument_id: InstrumentId,
326        message: &OKXOrderMsg,
327    ) -> LinkedChildResolution {
328        let candidate_client_order_ids = self
329            .order_contexts
330            .iter()
331            .filter_map(|entry| {
332                (entry.identity.instrument_id == instrument_id)
333                    .then_some(entry.identity.client_order_id)
334            })
335            .collect::<Vec<_>>();
336        let bindings = self.lifecycle_bindings.lock();
337        if let Some(client_order_id) = bindings.client_order_id(&parent_venue_order_id) {
338            return LinkedChildResolution::Bound(client_order_id);
339        }
340
341        let mut pending = self.pending_linked_children.lock();
342        let candidate_client_order_ids = candidate_client_order_ids
343            .iter()
344            .filter(|client_order_id| {
345                self.order_contexts.contains_key(client_order_id)
346                    && !bindings.venue_by_client.contains_key(client_order_id)
347            })
348            .copied()
349            .collect::<Vec<_>>();
350
351        if candidate_client_order_ids.is_empty() {
352            return LinkedChildResolution::External;
353        }
354
355        pending.push_back(PendingLinkedChild {
356            parent_venue_order_id,
357            candidate_client_order_ids,
358            message: message.clone(),
359        });
360        LinkedChildResolution::Held
361    }
362
363    fn take_routable_linked_children(&self) -> Vec<OKXOrderMsg> {
364        if self.pending_linked_children.lock().is_empty() {
365            return Vec::new();
366        }
367
368        let bindings = self.lifecycle_bindings.lock();
369        let mut pending = self.pending_linked_children.lock();
370        let mut held = VecDeque::with_capacity(pending.len());
371        let mut routable = Vec::new();
372
373        while let Some(child) = pending.pop_front() {
374            let is_bound = bindings
375                .client_order_id(&child.parent_venue_order_id)
376                .is_some();
377            let is_pending = child
378                .candidate_client_order_ids
379                .iter()
380                .any(|client_order_id| {
381                    self.order_contexts.contains_key(client_order_id)
382                        && !bindings.venue_by_client.contains_key(client_order_id)
383                });
384
385            if is_bound || !is_pending {
386                routable.push(child.message);
387            } else {
388                held.push_back(child);
389            }
390        }
391
392        *pending = held;
393        routable
394    }
395}
396
397impl WsDispatchState {
398    /// Returns whether acceptance was already emitted for the order.
399    #[must_use]
400    pub fn contains_accepted(&self, cid: &ClientOrderId) -> bool {
401        self.accepted_venue_order_ids.lock().contains_key(cid)
402    }
403
404    /// Records that acceptance was emitted for the order.
405    pub fn insert_accepted(&self, cid: ClientOrderId, venue_order_id: VenueOrderId) {
406        self.accepted_venue_order_ids
407            .lock()
408            .insert(cid, venue_order_id);
409    }
410
411    fn accepted_venue_order_id(&self, cid: &ClientOrderId) -> Option<VenueOrderId> {
412        self.accepted_venue_order_ids.lock().get(cid).copied()
413    }
414
415    /// Returns whether the order was already triggered.
416    #[must_use]
417    pub fn contains_triggered(&self, cid: &ClientOrderId) -> bool {
418        self.triggered_orders.contains(cid)
419    }
420
421    /// Records that the order was triggered.
422    pub fn insert_triggered(&self, cid: ClientOrderId) {
423        let _ = self.triggered_orders.insert(cid);
424    }
425
426    /// Returns whether the order was already filled.
427    #[must_use]
428    pub fn contains_filled(&self, cid: &ClientOrderId) -> bool {
429        self.filled_orders.contains(cid)
430    }
431
432    /// Records that the order was filled.
433    pub fn insert_filled(&self, cid: ClientOrderId) {
434        let _ = self.filled_orders.insert(cid);
435    }
436
437    /// Returns whether the order already reached a terminal state.
438    #[must_use]
439    pub fn contains_terminal(&self, cid: &ClientOrderId) -> bool {
440        self.terminal_orders.contains(cid)
441    }
442
443    /// Records that the order reached a terminal state.
444    pub fn insert_terminal(&self, cid: ClientOrderId) {
445        let _ = self.terminal_orders.insert(cid);
446    }
447
448    /// Returns `true` if this trade was already emitted (duplicate).
449    /// Uses atomic insert to avoid TOCTOU races between concurrent streams.
450    pub fn check_and_insert_trade(&self, trade_id: TradeId) -> bool {
451        !self.emitted_trades.insert(trade_id)
452    }
453
454    #[must_use]
455    pub fn contains_trade(&self, trade_id: &TradeId) -> bool {
456        self.emitted_trades.contains(trade_id)
457    }
458
459    fn remove_accepted(&self, cid: &ClientOrderId) {
460        self.accepted_venue_order_ids.lock().remove(cid);
461    }
462
463    fn remove_triggered(&self, cid: &ClientOrderId) {
464        self.triggered_orders.remove(cid);
465    }
466
467    fn remove_filled(&self, cid: &ClientOrderId) {
468        self.filled_orders.remove(cid);
469    }
470
471    fn insert_post_only_rejection(&self, order_id: Ustr) {
472        let _ = self.post_only_rejections.insert(order_id);
473    }
474
475    fn contains_post_only_rejection(&self, order_id: &Ustr) -> bool {
476        self.post_only_rejections.contains(order_id)
477    }
478}
479
480/// Dispatches a WebSocket message with cross-stream deduplication.
481///
482/// For orders with a tracked identity (submitted through this client), produces
483/// proper order events (`OrderAccepted`, `OrderCanceled`, `OrderFilled`, etc.).
484/// For untracked orders (external or pre-existing), falls back to execution
485/// reports for downstream reconciliation.
486#[expect(clippy::too_many_arguments)]
487pub fn dispatch_ws_message(
488    message: OKXWsMessage,
489    emitter: &ExecutionEventEmitter,
490    state: &WsDispatchState,
491    account_id: AccountId,
492    account_type: AccountType,
493    instruments: &AtomicMap<Ustr, InstrumentAny>,
494    fee_cache: &mut FeeCache,
495    filled_qty_cache: &mut FilledQtyCache,
496    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
497    clock: &AtomicTime,
498) {
499    let guard = instruments.load();
500    let instruments: &AHashMap<Ustr, InstrumentAny> = &guard;
501
502    match message {
503        OKXWsMessage::Orders(order_msgs) => {
504            let ts_init = clock.get_time_ns();
505            let pending_order_msgs = state.take_routable_linked_children();
506            if !pending_order_msgs.is_empty() {
507                dispatch_order_messages(
508                    &pending_order_msgs,
509                    emitter,
510                    state,
511                    account_id,
512                    instruments,
513                    fee_cache,
514                    filled_qty_cache,
515                    order_state_cache,
516                    ts_init,
517                );
518            }
519            dispatch_order_messages(
520                &order_msgs,
521                emitter,
522                state,
523                account_id,
524                instruments,
525                fee_cache,
526                filled_qty_cache,
527                order_state_cache,
528                ts_init,
529            );
530        }
531        OKXWsMessage::SpreadOrders(order_msgs) => {
532            let ts_init = clock.get_time_ns();
533            dispatch_spread_order_messages(
534                &order_msgs,
535                emitter,
536                state,
537                account_id,
538                instruments,
539                filled_qty_cache,
540                order_state_cache,
541                ts_init,
542            );
543        }
544        OKXWsMessage::AlgoOrders(algo_msgs) => {
545            let ts_init = clock.get_time_ns();
546            for msg in &algo_msgs {
547                dispatch_algo_order_message(msg, emitter, state, account_id, instruments, ts_init);
548            }
549        }
550        OKXWsMessage::Account(data) => {
551            let ts_init = clock.get_time_ns();
552
553            match serde_json::from_value::<Vec<OKXAccount>>(data) {
554                Ok(accounts) => {
555                    for account in &accounts {
556                        match crate::common::parse::parse_account_state(
557                            account,
558                            account_id,
559                            account_type,
560                            ts_init,
561                        ) {
562                            Ok(account_state) => emitter.send_account_state(account_state),
563                            Err(e) => log::error!("Failed to parse account state: {e}"),
564                        }
565                    }
566                }
567                Err(e) => log::error!("Failed to deserialize account data: {e}"),
568            }
569        }
570        OKXWsMessage::Positions(data) => {
571            let ts_init = clock.get_time_ns();
572
573            match serde_json::from_value::<Vec<OKXPosition>>(data) {
574                Ok(positions) => {
575                    for position in positions {
576                        let Some(instrument) = instruments.get(&position.inst_id) else {
577                            log::warn!("No cached instrument for position: {}", position.inst_id);
578                            continue;
579                        };
580                        let instrument_id = instrument.id();
581                        let size_precision = instrument.size_precision();
582
583                        match crate::common::parse::parse_position_status_report(
584                            &position,
585                            account_id,
586                            instrument_id,
587                            size_precision,
588                            ts_init,
589                        ) {
590                            Ok(report) => emitter.send_position_report(report),
591                            Err(e) => log::error!("Failed to parse position report: {e}"),
592                        }
593                    }
594                }
595                Err(e) => log::error!("Failed to deserialize positions data: {e}"),
596            }
597        }
598        OKXWsMessage::OrderResponse {
599            id,
600            op,
601            code,
602            msg,
603            data,
604        } => {
605            let ts_init = clock.get_time_ns();
606
607            for item in &data {
608                let s_code = item
609                    .get(OKX_FIELD_SCODE)
610                    .and_then(|v| v.as_str())
611                    .unwrap_or("");
612                let s_msg = item
613                    .get(OKX_FIELD_SMSG)
614                    .and_then(|v| v.as_str())
615                    .unwrap_or("");
616                let sub_code = item
617                    .get(OKX_FIELD_SUBCODE)
618                    .and_then(|v| v.as_str())
619                    .unwrap_or("");
620                let reason = format_order_response_reason(s_code, s_msg, sub_code);
621                let cl_ord_id = item
622                    .get(OKX_FIELD_CLORDID)
623                    .and_then(|v| v.as_str())
624                    .unwrap_or("");
625
626                if s_code == OKX_SUCCESS_CODE {
627                    log::debug!("Order response ok: op={op:?} cl_ord_id={cl_ord_id}");
628                    match op {
629                        OKXWsOperation::Order
630                        | OKXWsOperation::BatchOrders
631                        | OKXWsOperation::OrderAlgo => {
632                            state.pending_orders.remove(cl_ord_id);
633                        }
634                        OKXWsOperation::CancelOrder
635                        | OKXWsOperation::BatchCancelOrders
636                        | OKXWsOperation::MassCancel
637                        | OKXWsOperation::CancelAlgos => {
638                            state.pending_cancels.remove(cl_ord_id);
639                        }
640                        OKXWsOperation::AmendOrder | OKXWsOperation::BatchAmendOrders => {
641                            state.pending_amends.remove(cl_ord_id);
642                        }
643                        _ => {}
644                    }
645                    continue;
646                }
647
648                let Some(client_order_id) = parse_client_order_id(cl_ord_id) else {
649                    log::warn!(
650                        "Order response error without client_order_id: \
651                         op={op:?} s_code={s_code} s_msg={s_msg}"
652                    );
653                    continue;
654                };
655
656                let Some(ident) = state.order_identity(client_order_id) else {
657                    log::warn!(
658                        "Order response error for untracked order: \
659                         op={op:?} cl_ord_id={cl_ord_id} s_code={s_code} s_msg={s_msg}"
660                    );
661                    continue;
662                };
663
664                let venue_order_id = item
665                    .get("ordId")
666                    .and_then(|v| v.as_str())
667                    .filter(|s| !s.is_empty())
668                    .map(VenueOrderId::new);
669
670                match classify_okx_venue_code(s_code, reason.clone()) {
671                    CommandFailure::Ambiguous(reason) => {
672                        log::warn!(
673                            "Ambiguous order response for {client_order_id}, awaiting reconciliation: \
674                             op={op:?} s_code={s_code} {reason}"
675                        );
676                        continue;
677                    }
678                    CommandFailure::NotSent(_) => {
679                        log::warn!(
680                            "Unexpected NotSent classification for venue order response: \
681                             op={op:?} cl_ord_id={cl_ord_id} s_code={s_code}"
682                        );
683                        continue;
684                    }
685                    CommandFailure::VenueRejected(_) => {}
686                }
687
688                match op {
689                    OKXWsOperation::Order | OKXWsOperation::BatchOrders => {
690                        state.remove_order_tracking(client_order_id);
691                        state.pending_orders.remove(cl_ord_id);
692                        emitter.emit_order_rejected_event(
693                            ident.strategy_id,
694                            ident.instrument_id,
695                            client_order_id,
696                            &reason,
697                            ts_init,
698                            false,
699                        );
700                    }
701                    OKXWsOperation::CancelOrder
702                    | OKXWsOperation::BatchCancelOrders
703                    | OKXWsOperation::MassCancel => {
704                        state.pending_cancels.remove(cl_ord_id);
705                        emitter.emit_order_cancel_rejected_event(
706                            ident.strategy_id,
707                            ident.instrument_id,
708                            client_order_id,
709                            venue_order_id,
710                            &reason,
711                            ts_init,
712                        );
713                    }
714                    OKXWsOperation::AmendOrder | OKXWsOperation::BatchAmendOrders => {
715                        state.pending_amends.remove(cl_ord_id);
716                        emitter.emit_order_modify_rejected_event(
717                            ident.strategy_id,
718                            ident.instrument_id,
719                            client_order_id,
720                            venue_order_id,
721                            &reason,
722                            ts_init,
723                        );
724                    }
725                    _ => {
726                        log::warn!(
727                            "Order response error for unhandled op: \
728                             op={op:?} cl_ord_id={cl_ord_id} s_code={s_code} s_msg={s_msg}"
729                        );
730                    }
731                }
732            }
733
734            if code != "0" && data.is_empty() {
735                log::warn!(
736                    "Order response error (no data): id={id:?} op={op:?} code={code} msg={msg}"
737                );
738            }
739        }
740        OKXWsMessage::SendFailed {
741            request_id,
742            client_order_ids,
743            op,
744            error,
745        } => {
746            let failure = classify_okx_ws_failure(&error);
747            let is_ambiguous = matches!(failure, CommandFailure::Ambiguous(_));
748            log::warn!(
749                "WebSocket send failed without structured venue response: \
750                 request_id={request_id}, client_order_ids={client_order_ids:?}, \
751                 op={op:?}, {failure:?}"
752            );
753
754            for client_order_id in client_order_ids {
755                let key = client_order_id.as_str();
756
757                match op {
758                    Some(
759                        OKXWsOperation::Order
760                        | OKXWsOperation::BatchOrders
761                        | OKXWsOperation::OrderAlgo,
762                    ) => {
763                        if !is_ambiguous {
764                            state.pending_orders.remove(key);
765                        }
766                        emit_send_failed_submit(&failure, state, emitter, clock, client_order_id);
767                    }
768                    Some(
769                        OKXWsOperation::CancelOrder
770                        | OKXWsOperation::BatchCancelOrders
771                        | OKXWsOperation::MassCancel
772                        | OKXWsOperation::CancelAlgos,
773                    ) => {
774                        if !is_ambiguous {
775                            state.pending_cancels.remove(key);
776                        }
777                    }
778                    Some(OKXWsOperation::AmendOrder | OKXWsOperation::BatchAmendOrders) => {
779                        if !is_ambiguous {
780                            state.pending_amends.remove(key);
781                        }
782                        emit_send_failed_modify(&failure, state, emitter, clock, client_order_id);
783                    }
784                    _ => {}
785                }
786            }
787        }
788        OKXWsMessage::ChannelData { channel, .. } => {
789            log::debug!("Ignoring data channel message on execution client: {channel:?}");
790        }
791        OKXWsMessage::SubscriptionFailed {
792            channel,
793            inst_id,
794            code,
795            msg,
796        } => {
797            log::error!(
798                "OKX rejected {channel:?} subscription for {inst_id:?} \
799                 (code={code}, msg={msg}); execution updates for it will not flow"
800            );
801        }
802        OKXWsMessage::LiquidationWarnings(warnings) => {
803            for warning in warnings {
804                log::warn!(
805                    "Liquidation warning: inst_id={}, pos_side={:?}, pos={}, mgn_ratio={}, mark_px={}, mgn_mode={:?}",
806                    warning.inst_id,
807                    warning.pos_side,
808                    warning.pos,
809                    warning.mgn_ratio,
810                    warning.mark_px,
811                    warning.mgn_mode,
812                );
813            }
814        }
815        OKXWsMessage::BookData { .. }
816        | OKXWsMessage::RpiBookData { .. }
817        | OKXWsMessage::Instruments(_) => {
818            log::debug!("Ignoring data message on execution client");
819        }
820        OKXWsMessage::Error(e) => {
821            log::warn!(
822                "Websocket error: code={} message={} conn_id={:?}",
823                e.code,
824                e.message,
825                e.conn_id
826            );
827        }
828        OKXWsMessage::Reconnected => {
829            log::info!("Websocket reconnected");
830        }
831        OKXWsMessage::Authenticated => {
832            log::debug!("Websocket authenticated");
833        }
834    }
835}
836
837fn route_algo_order_message(
838    msg: &OKXAlgoOrderMsg,
839    state: &WsDispatchState,
840) -> ExecutionUpdateRoute {
841    if matches!(
842        msg.ord_type,
843        OKXAlgoOrderType::Iceberg
844            | OKXAlgoOrderType::SmartIceberg
845            | OKXAlgoOrderType::Twap
846            | OKXAlgoOrderType::Chase
847            | OKXAlgoOrderType::Other
848    ) || msg.state == OKXAlgoOrderStatus::Unknown
849    {
850        return ExecutionUpdateRoute::Suppressed;
851    }
852
853    let direct_client_order_id = parse_client_order_id(&msg.algo_cl_ord_id)
854        .or_else(|| parse_client_order_id(&msg.cl_ord_id));
855
856    if let Some(client_order_id) = direct_client_order_id {
857        if state.contains_terminal(&client_order_id) {
858            return ExecutionUpdateRoute::Suppressed;
859        }
860
861        if let Some(context) = state
862            .order_contexts
863            .get(&client_order_id)
864            .map(|entry| *entry)
865        {
866            return ExecutionUpdateRoute::Tracked(client_order_id, context);
867        }
868
869        if state.order_identities.contains_key(&client_order_id) {
870            return ExecutionUpdateRoute::Suppressed;
871        }
872    }
873
874    let parent_venue_order_id = VenueOrderId::new(msg.algo_id.as_str());
875    let client_order_id = {
876        let bindings = state.lifecycle_bindings.lock();
877        bindings.client_order_id(&parent_venue_order_id)
878    };
879
880    let Some(client_order_id) = client_order_id else {
881        return ExecutionUpdateRoute::External;
882    };
883
884    if state.contains_terminal(&client_order_id) {
885        return ExecutionUpdateRoute::Suppressed;
886    }
887
888    state
889        .order_contexts
890        .get(&client_order_id)
891        .map_or(ExecutionUpdateRoute::Suppressed, |entry| {
892            ExecutionUpdateRoute::Tracked(client_order_id, *entry)
893        })
894}
895
896fn dispatch_algo_order_message(
897    msg: &OKXAlgoOrderMsg,
898    emitter: &ExecutionEventEmitter,
899    state: &WsDispatchState,
900    account_id: AccountId,
901    instruments: &AHashMap<Ustr, InstrumentAny>,
902    ts_init: UnixNanos,
903) {
904    let route = route_algo_order_message(msg, state);
905
906    match route {
907        ExecutionUpdateRoute::External => {
908            match parse_algo_order_msg(msg, account_id, instruments, ts_init) {
909                Ok(Some(report)) => dispatch_execution_reports(vec![report], emitter, state),
910                Ok(None) => {}
911                Err(e) => log::error!("Failed to parse external algo order message: {e}"),
912            }
913        }
914        ExecutionUpdateRoute::Suppressed => {
915            log::debug!(
916                "Suppressing algo order update: algo_id={} state={:?}",
917                msg.algo_id,
918                msg.state,
919            );
920        }
921        ExecutionUpdateRoute::Tracked(client_order_id, context) => {
922            let Some(instrument) = instruments.get(&msg.inst_id) else {
923                log::warn!(
924                    "No instrument for {}, skipping algo order message",
925                    msg.inst_id
926                );
927                return;
928            };
929            dispatch_tracked_algo_order_message(
930                msg,
931                client_order_id,
932                context,
933                instrument,
934                emitter,
935                state,
936                account_id,
937                ts_init,
938            );
939        }
940    }
941}
942
943#[expect(
944    clippy::too_many_arguments,
945    reason = "tracked routing requires the resolved context, venue state, and event timestamps"
946)]
947fn dispatch_tracked_algo_order_message(
948    msg: &OKXAlgoOrderMsg,
949    client_order_id: ClientOrderId,
950    mut context: OrderContext,
951    instrument: &InstrumentAny,
952    emitter: &ExecutionEventEmitter,
953    state: &WsDispatchState,
954    account_id: AccountId,
955    ts_init: UnixNanos,
956) {
957    let parent_venue_order_id = VenueOrderId::new(msg.algo_id.as_str());
958    let ts_event = parse_millisecond_timestamp(msg.u_time);
959    let mut bindings = state.lifecycle_bindings.lock();
960    let mut is_terminal = false;
961    let mut binding = bindings
962        .venue_by_client
963        .get(&client_order_id)
964        .copied()
965        .unwrap_or(OrderVenueBinding {
966            parent: parent_venue_order_id,
967            child: None,
968        });
969
970    if binding.parent != parent_venue_order_id {
971        log::error!(
972            "Suppressing conflicting algo parent binding for {client_order_id}: expected={} received={parent_venue_order_id}",
973            binding.parent,
974        );
975        return;
976    }
977
978    if binding.child.is_none() {
979        context = refresh_algo_order_context(msg, context, instrument, account_id, ts_init);
980        state.track_order_context(context);
981    }
982
983    bindings
984        .client_by_parent
985        .insert(parent_venue_order_id, client_order_id);
986
987    match msg.state {
988        OKXAlgoOrderStatus::Live | OKXAlgoOrderStatus::Pause => {
989            if binding.child.is_some() {
990                log::debug!(
991                    "Suppressing stale algo parent acceptance for {client_order_id}: algo_id={}",
992                    msg.algo_id,
993                );
994            } else {
995                ensure_accepted_emitted(
996                    client_order_id,
997                    account_id,
998                    parent_venue_order_id,
999                    &context.identity,
1000                    emitter,
1001                    state,
1002                    ts_event,
1003                    ts_init,
1004                );
1005            }
1006        }
1007        OKXAlgoOrderStatus::Effective
1008        | OKXAlgoOrderStatus::OrderPlaced
1009        | OKXAlgoOrderStatus::PartiallyEffective
1010        | OKXAlgoOrderStatus::Filled
1011        | OKXAlgoOrderStatus::PartiallyFailed => {
1012            let child_venue_order_id = algo_child_venue_order_id(msg);
1013            if let Some(child_venue_order_id) = child_venue_order_id {
1014                bind_algo_child_and_emit_transition(
1015                    client_order_id,
1016                    parent_venue_order_id,
1017                    child_venue_order_id,
1018                    &mut binding,
1019                    context,
1020                    account_id,
1021                    ts_event,
1022                    ts_init,
1023                    emitter,
1024                    state,
1025                );
1026            } else {
1027                ensure_accepted_emitted(
1028                    client_order_id,
1029                    account_id,
1030                    parent_venue_order_id,
1031                    &context.identity,
1032                    emitter,
1033                    state,
1034                    ts_event,
1035                    ts_init,
1036                );
1037            }
1038
1039            if matches!(
1040                msg.state,
1041                OKXAlgoOrderStatus::Filled | OKXAlgoOrderStatus::PartiallyFailed
1042            ) {
1043                log::debug!(
1044                    "Deferring tracked algo {:?} update for {client_order_id} to regular child execution updates",
1045                    msg.state,
1046                );
1047            }
1048        }
1049        OKXAlgoOrderStatus::Canceled => {
1050            if binding.child.is_some() {
1051                log::debug!(
1052                    "Suppressing stale canceled algo parent for triggered order {client_order_id}"
1053                );
1054            } else {
1055                ensure_accepted_emitted(
1056                    client_order_id,
1057                    account_id,
1058                    parent_venue_order_id,
1059                    &context.identity,
1060                    emitter,
1061                    state,
1062                    ts_event,
1063                    ts_init,
1064                );
1065                let canceled = OrderCanceled::new(
1066                    emitter.trader_id(),
1067                    context.identity.strategy_id,
1068                    context.identity.instrument_id,
1069                    client_order_id,
1070                    UUID4::new(),
1071                    ts_event,
1072                    ts_init,
1073                    false,
1074                    Some(parent_venue_order_id),
1075                    Some(account_id),
1076                    None,
1077                );
1078                state.insert_terminal(client_order_id);
1079                state.remove_accepted(&client_order_id);
1080                state.remove_order_tracking(client_order_id);
1081                is_terminal = true;
1082                emitter.send_order_event(OrderEventAny::Canceled(canceled));
1083            }
1084        }
1085        OKXAlgoOrderStatus::OrderFailed => {
1086            if binding.child.is_some() {
1087                log::debug!(
1088                    "Suppressing stale failed algo parent for triggered order {client_order_id}"
1089                );
1090            } else {
1091                let reason = if msg.fail_code.is_empty() {
1092                    "OKX algo order failed"
1093                } else {
1094                    msg.fail_code.as_str()
1095                };
1096                let rejected = OrderRejected::new(
1097                    emitter.trader_id(),
1098                    context.identity.strategy_id,
1099                    context.identity.instrument_id,
1100                    client_order_id,
1101                    account_id,
1102                    Ustr::from(reason),
1103                    UUID4::new(),
1104                    ts_event,
1105                    ts_init,
1106                    false,
1107                    false,
1108                );
1109                state.insert_terminal(client_order_id);
1110                state.remove_accepted(&client_order_id);
1111                state.remove_order_tracking(client_order_id);
1112                is_terminal = true;
1113                emitter.send_order_event(OrderEventAny::Rejected(rejected));
1114            }
1115        }
1116        OKXAlgoOrderStatus::Unknown => {}
1117    }
1118
1119    if is_terminal {
1120        bindings.finish(client_order_id, binding);
1121    } else {
1122        bindings.venue_by_client.insert(client_order_id, binding);
1123    }
1124
1125    drop(bindings);
1126    state.linked_child_notify.notify_one();
1127}
1128
1129fn algo_child_venue_order_id(msg: &OKXAlgoOrderMsg) -> Option<VenueOrderId> {
1130    if !msg.ord_id.is_empty() {
1131        return Some(VenueOrderId::new(msg.ord_id.as_str()));
1132    }
1133
1134    match msg.ord_id_list.as_slice() {
1135        [order_id] if !order_id.is_empty() => Some(VenueOrderId::new(order_id.as_str())),
1136        [] => None,
1137        order_ids => {
1138            log::warn!(
1139                "Cannot bind algo order {} to {} triggered child IDs",
1140                msg.algo_id,
1141                order_ids.len(),
1142            );
1143            None
1144        }
1145    }
1146}
1147
1148fn refresh_algo_order_context(
1149    msg: &OKXAlgoOrderMsg,
1150    mut context: OrderContext,
1151    instrument: &InstrumentAny,
1152    account_id: AccountId,
1153    ts_init: UnixNanos,
1154) -> OrderContext {
1155    match parse_algo_order_status_report(msg, instrument, account_id, ts_init) {
1156        Ok(report) => {
1157            if !msg.sz.is_empty() {
1158                context.quantity = report.quantity;
1159            }
1160
1161            context.price = report.price;
1162            context.trigger_price = report.trigger_price;
1163            context.trigger_type = report.trigger_type;
1164        }
1165        Err(e) => {
1166            log::error!(
1167                "Failed to refresh tracked algo order context for {}: {e}",
1168                context.identity.client_order_id,
1169            );
1170            return context;
1171        }
1172    }
1173
1174    if !msg.actual_sz.is_empty() && msg.actual_sz != "0" {
1175        match parse_quantity(msg.actual_sz.as_str(), instrument.size_precision()) {
1176            Ok(quantity) => context.quantity = quantity,
1177            Err(e) => log::error!(
1178                "Failed to refresh tracked algo actual quantity for {}: {e}",
1179                context.identity.client_order_id,
1180            ),
1181        }
1182    }
1183
1184    context
1185}
1186
1187fn refresh_regular_child_context(
1188    msg: &OKXOrderMsg,
1189    mut context: OrderContext,
1190    instrument: &InstrumentAny,
1191) -> OrderContext {
1192    match parse_quantity(&msg.sz, instrument.size_precision()) {
1193        Ok(quantity) => context.quantity = quantity,
1194        Err(e) => log::error!(
1195            "Failed to refresh tracked child quantity for {}: {e}",
1196            context.identity.client_order_id,
1197        ),
1198    }
1199
1200    context.price = if is_market_price(&msg.px) {
1201        None
1202    } else {
1203        match parse_price(&msg.px, instrument.price_precision()) {
1204            Ok(price) => Some(price),
1205            Err(e) => {
1206                log::error!(
1207                    "Failed to refresh tracked child price for {}: {e}",
1208                    context.identity.client_order_id,
1209                );
1210                context.price
1211            }
1212        }
1213    };
1214    context
1215}
1216
1217#[expect(clippy::too_many_arguments)]
1218fn bind_algo_child_and_emit_transition(
1219    client_order_id: ClientOrderId,
1220    parent_venue_order_id: VenueOrderId,
1221    child_venue_order_id: VenueOrderId,
1222    binding: &mut OrderVenueBinding,
1223    context: OrderContext,
1224    account_id: AccountId,
1225    ts_event: UnixNanos,
1226    ts_init: UnixNanos,
1227    emitter: &ExecutionEventEmitter,
1228    state: &WsDispatchState,
1229) {
1230    if let Some(bound_child) = binding.child {
1231        if bound_child != child_venue_order_id {
1232            log::error!(
1233                "Suppressing conflicting algo child binding for {client_order_id}: expected={bound_child} received={child_venue_order_id}"
1234            );
1235        }
1236        return;
1237    }
1238
1239    binding.child = Some(child_venue_order_id);
1240    state.track_order_context(context);
1241    ensure_accepted_emitted(
1242        client_order_id,
1243        account_id,
1244        parent_venue_order_id,
1245        &context.identity,
1246        emitter,
1247        state,
1248        ts_event,
1249        ts_init,
1250    );
1251
1252    if state.accepted_venue_order_id(&client_order_id) != Some(child_venue_order_id) {
1253        state.insert_accepted(client_order_id, child_venue_order_id);
1254        emit_child_update(
1255            client_order_id,
1256            child_venue_order_id,
1257            context,
1258            account_id,
1259            ts_event,
1260            ts_init,
1261            emitter,
1262        );
1263    }
1264
1265    if !state.contains_triggered(&client_order_id) {
1266        state.insert_triggered(client_order_id);
1267
1268        if TRIGGERABLE_ORDER_TYPES.contains(&context.identity.order_type) {
1269            let triggered = OrderTriggered::new(
1270                emitter.trader_id(),
1271                context.identity.strategy_id,
1272                context.identity.instrument_id,
1273                client_order_id,
1274                UUID4::new(),
1275                ts_event,
1276                ts_init,
1277                false,
1278                Some(child_venue_order_id),
1279                Some(account_id),
1280            );
1281            emitter.send_order_event(OrderEventAny::Triggered(triggered));
1282        }
1283    }
1284}
1285
1286#[expect(clippy::too_many_arguments)]
1287fn refresh_bound_child_and_emit_update(
1288    client_order_id: ClientOrderId,
1289    child_venue_order_id: VenueOrderId,
1290    context: OrderContext,
1291    account_id: AccountId,
1292    ts_event: UnixNanos,
1293    ts_init: UnixNanos,
1294    emitter: &ExecutionEventEmitter,
1295    state: &WsDispatchState,
1296    emit_update: bool,
1297) {
1298    let terms_changed = state
1299        .order_contexts
1300        .get(&client_order_id)
1301        .is_some_and(|previous| {
1302            previous.quantity != context.quantity
1303                || previous.price != context.price
1304                || previous.trigger_price != context.trigger_price
1305        });
1306    state.track_order_context(context);
1307
1308    if !terms_changed || !emit_update {
1309        return;
1310    }
1311
1312    state.insert_accepted(client_order_id, child_venue_order_id);
1313    emit_child_update(
1314        client_order_id,
1315        child_venue_order_id,
1316        context,
1317        account_id,
1318        ts_event,
1319        ts_init,
1320        emitter,
1321    );
1322}
1323
1324fn emit_child_update(
1325    client_order_id: ClientOrderId,
1326    child_venue_order_id: VenueOrderId,
1327    context: OrderContext,
1328    account_id: AccountId,
1329    ts_event: UnixNanos,
1330    ts_init: UnixNanos,
1331    emitter: &ExecutionEventEmitter,
1332) {
1333    let updated = OrderUpdated::new(
1334        emitter.trader_id(),
1335        context.identity.strategy_id,
1336        context.identity.instrument_id,
1337        client_order_id,
1338        context.quantity,
1339        UUID4::new(),
1340        ts_event,
1341        ts_init,
1342        false,
1343        Some(child_venue_order_id),
1344        Some(account_id),
1345        context.price,
1346        context.trigger_price,
1347        None,
1348        false,
1349    );
1350    emitter.send_order_event(OrderEventAny::Updated(updated));
1351}
1352
1353/// Dispatches order messages, producing proper order events for tracked orders
1354/// and falling back to execution reports for untracked/external orders.
1355#[expect(clippy::too_many_arguments)]
1356fn dispatch_order_messages(
1357    order_msgs: &[OKXOrderMsg],
1358    emitter: &ExecutionEventEmitter,
1359    state: &WsDispatchState,
1360    account_id: AccountId,
1361    instruments: &AHashMap<Ustr, InstrumentAny>,
1362    fee_cache: &mut FeeCache,
1363    filled_qty_cache: &mut FilledQtyCache,
1364    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
1365    ts_init: UnixNanos,
1366) {
1367    for msg in order_msgs {
1368        let Some(instrument) = instruments.get(&msg.inst_id) else {
1369            log::warn!("No instrument for {}, skipping order message", msg.inst_id);
1370            continue;
1371        };
1372
1373        let direct_client_order_id = parse_client_order_id(&msg.cl_ord_id);
1374        let parent_client_order_id = msg
1375            .algo_cl_ord_id
1376            .as_deref()
1377            .and_then(parse_client_order_id);
1378        let linked_parent_venue_order_id = msg
1379            .algo_id
1380            .as_deref()
1381            .filter(|value| !value.is_empty())
1382            .or_else(|| {
1383                msg.linked_algo_ord
1384                    .as_ref()
1385                    .map(|linked| linked.algo_id.as_str())
1386                    .filter(|value| !value.is_empty())
1387            })
1388            .map(VenueOrderId::new);
1389
1390        // Triggered child orders may have a generated or empty cl_ord_id.
1391        // Resolve the tracked parent before falling back to a report.
1392        let direct_resolution = [direct_client_order_id, parent_client_order_id]
1393            .into_iter()
1394            .flatten()
1395            .find_map(|client_order_id| {
1396                state
1397                    .order_identity(client_order_id)
1398                    .map(|identity| (client_order_id, Some(identity)))
1399            });
1400
1401        let bound_client_order_id = if direct_resolution.is_none()
1402            && let Some(parent_venue_order_id) = linked_parent_venue_order_id
1403        {
1404            match state.resolve_or_hold_linked_child(parent_venue_order_id, instrument.id(), msg) {
1405                LinkedChildResolution::Bound(client_order_id) => Some(client_order_id),
1406                LinkedChildResolution::Held => {
1407                    log::debug!(
1408                        "Holding linked child update during algo parent binding: ord_id={} parent_id={parent_venue_order_id}",
1409                        msg.ord_id,
1410                    );
1411                    continue;
1412                }
1413                LinkedChildResolution::External => None,
1414            }
1415        } else {
1416            None
1417        };
1418
1419        let resolved = direct_resolution
1420            .or_else(|| {
1421                bound_client_order_id.and_then(|client_order_id| {
1422                    state
1423                        .order_identity(client_order_id)
1424                        .map(|identity| (client_order_id, Some(identity)))
1425                })
1426            })
1427            .or_else(|| {
1428                bound_client_order_id
1429                    .or(parent_client_order_id)
1430                    .or(direct_client_order_id)
1431                    .map(|client_order_id| (client_order_id, None))
1432            });
1433
1434        let Some((client_order_id, identity)) = resolved else {
1435            log::debug!(
1436                "Order without client or algo client order ID (ord_id={}), sending as report",
1437                msg.ord_id
1438            );
1439            dispatch_order_msg_as_report(
1440                msg,
1441                account_id,
1442                instruments,
1443                fee_cache,
1444                filled_qty_cache,
1445                emitter,
1446                state,
1447                ts_init,
1448            );
1449            continue;
1450        };
1451
1452        if let Some(ident) = identity {
1453            let context = state
1454                .order_contexts
1455                .get(&client_order_id)
1456                .map(|entry| refresh_regular_child_context(msg, *entry, instrument));
1457            let mut lifecycle_bindings = context.map(|_| state.lifecycle_bindings.lock());
1458
1459            if let (Some(context), Some(bindings)) = (context, lifecycle_bindings.as_mut()) {
1460                let parent_venue_order_id = linked_parent_venue_order_id.or_else(|| {
1461                    bindings
1462                        .venue_by_client
1463                        .get(&client_order_id)
1464                        .map(|binding| binding.parent)
1465                });
1466
1467                let child_venue_order_id = VenueOrderId::new(msg.ord_id);
1468                let parent_venue_order_id = parent_venue_order_id.unwrap_or(child_venue_order_id);
1469                let mut binding = bindings
1470                    .venue_by_client
1471                    .get(&client_order_id)
1472                    .copied()
1473                    .unwrap_or(OrderVenueBinding {
1474                        parent: parent_venue_order_id,
1475                        child: None,
1476                    });
1477
1478                if binding.parent != parent_venue_order_id {
1479                    log::error!(
1480                        "Suppressing conflicting child parent binding for {client_order_id}: expected={} received={parent_venue_order_id}",
1481                        binding.parent,
1482                    );
1483                    continue;
1484                }
1485
1486                bindings
1487                    .client_by_parent
1488                    .insert(parent_venue_order_id, client_order_id);
1489                let ts_event = parse_millisecond_timestamp(msg.u_time);
1490
1491                if binding.child == Some(child_venue_order_id) {
1492                    refresh_bound_child_and_emit_update(
1493                        client_order_id,
1494                        child_venue_order_id,
1495                        context,
1496                        account_id,
1497                        ts_event,
1498                        ts_init,
1499                        emitter,
1500                        state,
1501                        !order_state_cache.contains_key(&client_order_id),
1502                    );
1503                } else {
1504                    bind_algo_child_and_emit_transition(
1505                        client_order_id,
1506                        parent_venue_order_id,
1507                        child_venue_order_id,
1508                        &mut binding,
1509                        context,
1510                        account_id,
1511                        ts_event,
1512                        ts_init,
1513                        emitter,
1514                        state,
1515                    );
1516                }
1517
1518                bindings.venue_by_client.insert(client_order_id, binding);
1519            }
1520
1521            let is_post_only_cancel = is_post_only_auto_cancel(msg);
1522
1523            if is_post_only_cancel
1524                || (!state.contains_accepted(&client_order_id) && is_unfilled_rpi_cancel(msg))
1525            {
1526                if is_post_only_cancel {
1527                    state.insert_post_only_rejection(msg.ord_id);
1528                }
1529
1530                let ts_event = parse_millisecond_timestamp(msg.u_time);
1531                let reason = if msg.ord_type == OKXOrderType::Rpi {
1532                    msg.cancel_source_reason
1533                        .as_deref()
1534                        .filter(|reason| !reason.is_empty())
1535                        .unwrap_or("RPI order canceled before acceptance")
1536                } else {
1537                    "Post-only order would have taken liquidity"
1538                };
1539                let rejected = OrderRejected::new(
1540                    emitter.trader_id(),
1541                    ident.strategy_id,
1542                    instrument.id(),
1543                    client_order_id,
1544                    account_id,
1545                    Ustr::from(reason),
1546                    UUID4::new(),
1547                    ts_event,
1548                    ts_init,
1549                    false,
1550                    true, // due_post_only
1551                );
1552                state.remove_order_tracking(client_order_id);
1553                if let Some(bindings) = lifecycle_bindings.as_mut() {
1554                    state.insert_terminal(client_order_id);
1555                    if let Some(binding) = bindings.venue_by_client.get(&client_order_id).copied() {
1556                        bindings.finish(client_order_id, binding);
1557                    }
1558                }
1559
1560                order_state_cache.remove(&client_order_id);
1561                fee_cache.remove(&msg.ord_id);
1562                filled_qty_cache.remove(&msg.ord_id);
1563                emitter.send_order_event(OrderEventAny::Rejected(rejected));
1564                continue;
1565            }
1566
1567            let previous_fee = fee_cache.get(&msg.ord_id).copied();
1568            let previous_filled_qty = filled_qty_cache.get(&msg.ord_id).copied();
1569            let previous_state = order_state_cache.get(&client_order_id);
1570
1571            match parse_order_event(
1572                msg,
1573                client_order_id,
1574                account_id,
1575                emitter.trader_id(),
1576                ident.strategy_id,
1577                instrument,
1578                previous_fee,
1579                previous_filled_qty,
1580                previous_state,
1581                ts_init,
1582            ) {
1583                Ok(event) => {
1584                    update_order_state_cache(msg, instrument, client_order_id, order_state_cache);
1585                    dispatch_parsed_order_event(
1586                        event,
1587                        client_order_id,
1588                        account_id,
1589                        VenueOrderId::new(msg.ord_id),
1590                        &ident,
1591                        instrument,
1592                        msg.state,
1593                        emitter,
1594                        state,
1595                        order_state_cache,
1596                        ts_init,
1597                    );
1598
1599                    if state.contains_terminal(&client_order_id)
1600                        && let Some(bindings) = lifecycle_bindings.as_mut()
1601                        && let Some(binding) =
1602                            bindings.venue_by_client.get(&client_order_id).copied()
1603                    {
1604                        bindings.finish(client_order_id, binding);
1605                    }
1606
1607                    update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
1608                }
1609                Err(e) => log::error!("Failed to parse order event for {client_order_id}: {e}"),
1610            }
1611        } else if state.contains_terminal(&client_order_id) {
1612            dispatch_terminal_order_fill_as_report(
1613                msg,
1614                client_order_id,
1615                account_id,
1616                instruments,
1617                fee_cache,
1618                filled_qty_cache,
1619                emitter,
1620                state,
1621                ts_init,
1622            );
1623        } else if is_post_only_auto_cancel(msg) && state.contains_post_only_rejection(&msg.ord_id) {
1624            log::debug!(
1625                "Skipping replayed post-only rejection for {client_order_id}: ord_id={}",
1626                msg.ord_id
1627            );
1628        } else {
1629            log::debug!(
1630                "Untracked order {client_order_id} (ord_id={}), sending as report for reconciliation",
1631                msg.ord_id
1632            );
1633            dispatch_order_msg_as_report(
1634                msg,
1635                account_id,
1636                instruments,
1637                fee_cache,
1638                filled_qty_cache,
1639                emitter,
1640                state,
1641                ts_init,
1642            );
1643        }
1644    }
1645}
1646
1647#[expect(clippy::too_many_arguments)]
1648fn dispatch_spread_order_messages(
1649    order_msgs: &[OKXSpreadOrder],
1650    emitter: &ExecutionEventEmitter,
1651    state: &WsDispatchState,
1652    account_id: AccountId,
1653    instruments: &AHashMap<Ustr, InstrumentAny>,
1654    filled_qty_cache: &mut FilledQtyCache,
1655    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
1656    ts_init: UnixNanos,
1657) {
1658    for msg in order_msgs {
1659        let Some(instrument) = instruments.get(&msg.sprd_id) else {
1660            log::warn!(
1661                "No instrument for {}, skipping spread order message",
1662                msg.sprd_id
1663            );
1664            continue;
1665        };
1666
1667        let Some(client_order_id) = parse_client_order_id(msg.cl_ord_id.as_str()) else {
1668            log::debug!(
1669                "Spread order without client_order_id (ord_id={}), sending as report",
1670                msg.ord_id
1671            );
1672            dispatch_spread_order_msg_as_report(
1673                msg,
1674                account_id,
1675                instruments,
1676                filled_qty_cache,
1677                emitter,
1678                state,
1679                ts_init,
1680            );
1681            continue;
1682        };
1683
1684        let identity = state.order_identity(client_order_id);
1685
1686        if let Some(ident) = identity {
1687            if is_spread_post_only_auto_cancel(msg) {
1688                let ts_event = msg
1689                    .u_time
1690                    .or(msg.c_time)
1691                    .map_or(ts_init, parse_millisecond_timestamp);
1692                let rejected = OrderRejected::new(
1693                    emitter.trader_id(),
1694                    ident.strategy_id,
1695                    instrument.id(),
1696                    client_order_id,
1697                    account_id,
1698                    Ustr::from(OKX_POST_ONLY_CANCEL_REASON),
1699                    UUID4::new(),
1700                    ts_event,
1701                    ts_init,
1702                    false,
1703                    true,
1704                );
1705                state.remove_order_tracking(client_order_id);
1706                order_state_cache.remove(&client_order_id);
1707                filled_qty_cache.remove(&msg.ord_id);
1708                emitter.send_order_event(OrderEventAny::Rejected(rejected));
1709                continue;
1710            }
1711
1712            let previous_filled_qty = filled_qty_cache.get(&msg.ord_id).copied();
1713            let previous_state = order_state_cache.get(&client_order_id);
1714
1715            match parse_spread_order_event(
1716                msg,
1717                client_order_id,
1718                account_id,
1719                emitter.trader_id(),
1720                ident.strategy_id,
1721                instrument,
1722                previous_filled_qty,
1723                previous_state,
1724                ts_init,
1725            ) {
1726                Ok(event) => {
1727                    update_spread_order_state_cache(
1728                        msg,
1729                        instrument,
1730                        client_order_id,
1731                        order_state_cache,
1732                    );
1733                    dispatch_parsed_order_event(
1734                        event,
1735                        client_order_id,
1736                        account_id,
1737                        VenueOrderId::new(msg.ord_id.as_str()),
1738                        &ident,
1739                        instrument,
1740                        msg.state,
1741                        emitter,
1742                        state,
1743                        order_state_cache,
1744                        ts_init,
1745                    );
1746                    update_spread_fill_cache(msg, instrument, filled_qty_cache);
1747                }
1748                Err(e) => {
1749                    log::error!("Failed to parse spread order event for {client_order_id}: {e}");
1750                }
1751            }
1752        } else {
1753            log::debug!(
1754                "Untracked spread order {client_order_id} (ord_id={}), sending as report for reconciliation",
1755                msg.ord_id
1756            );
1757            dispatch_spread_order_msg_as_report(
1758                msg,
1759                account_id,
1760                instruments,
1761                filled_qty_cache,
1762                emitter,
1763                state,
1764                ts_init,
1765            );
1766        }
1767    }
1768}
1769
1770/// Dispatches a parsed order event as a proper `OrderEventAny`.
1771///
1772/// Guarantees the `Submitted -> Accepted -> ...` lifecycle by synthesizing
1773/// `OrderAccepted` before any other event when one has not yet been emitted.
1774/// Duplicate `Accepted` events (e.g. from reconnect replays) are suppressed.
1775#[expect(clippy::too_many_arguments)]
1776fn dispatch_parsed_order_event(
1777    event: ParsedOrderEvent,
1778    client_order_id: ClientOrderId,
1779    account_id: AccountId,
1780    venue_order_id: VenueOrderId,
1781    identity: &OrderIdentity,
1782    instrument: &InstrumentAny,
1783    venue_status: OKXOrderStatus,
1784    emitter: &ExecutionEventEmitter,
1785    state: &WsDispatchState,
1786    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
1787    ts_init: UnixNanos,
1788) {
1789    let is_terminal;
1790
1791    match event {
1792        ParsedOrderEvent::Accepted(e) => {
1793            if state.contains_filled(&client_order_id) || state.contains_terminal(&client_order_id)
1794            {
1795                log::debug!("Skipping duplicate Accepted for {client_order_id}");
1796                return;
1797            }
1798
1799            if state.contains_accepted(&client_order_id) {
1800                emit_venue_order_id_update_if_changed(
1801                    client_order_id,
1802                    account_id,
1803                    venue_order_id,
1804                    identity,
1805                    e.ts_event,
1806                    emitter,
1807                    state,
1808                    order_state_cache,
1809                    ts_init,
1810                );
1811                return;
1812            }
1813
1814            if state.contains_triggered(&client_order_id) {
1815                log::debug!("Skipping duplicate Accepted for {client_order_id}");
1816                return;
1817            }
1818
1819            state.insert_accepted(client_order_id, venue_order_id);
1820            is_terminal = false;
1821            emitter.send_order_event(OrderEventAny::Accepted(e));
1822        }
1823        ParsedOrderEvent::Triggered(e) => {
1824            if state.contains_filled(&client_order_id) {
1825                log::debug!("Skipping stale Triggered for {client_order_id} (already filled)");
1826                return;
1827            }
1828
1829            if !TRIGGERABLE_ORDER_TYPES.contains(&identity.order_type) {
1830                log::debug!(
1831                    "Skipping OrderTriggered for {} order {client_order_id}: market-style stops have no TRIGGERED state",
1832                    identity.order_type,
1833                );
1834                state.insert_triggered(client_order_id);
1835                return;
1836            }
1837
1838            ensure_accepted_emitted(
1839                client_order_id,
1840                account_id,
1841                venue_order_id,
1842                identity,
1843                emitter,
1844                state,
1845                ts_init,
1846                ts_init,
1847            );
1848            state.insert_triggered(client_order_id);
1849            is_terminal = false;
1850            emitter.send_order_event(OrderEventAny::Triggered(e));
1851        }
1852        ParsedOrderEvent::Canceled(e) => {
1853            ensure_accepted_emitted(
1854                client_order_id,
1855                account_id,
1856                venue_order_id,
1857                identity,
1858                emitter,
1859                state,
1860                ts_init,
1861                ts_init,
1862            );
1863            state.remove_triggered(&client_order_id);
1864            state.remove_filled(&client_order_id);
1865            is_terminal = true;
1866            emitter.send_order_event(OrderEventAny::Canceled(e));
1867        }
1868        ParsedOrderEvent::Expired(e) => {
1869            ensure_accepted_emitted(
1870                client_order_id,
1871                account_id,
1872                venue_order_id,
1873                identity,
1874                emitter,
1875                state,
1876                ts_init,
1877                ts_init,
1878            );
1879            state.remove_triggered(&client_order_id);
1880            state.remove_filled(&client_order_id);
1881            is_terminal = true;
1882            emitter.send_order_event(OrderEventAny::Expired(e));
1883        }
1884        ParsedOrderEvent::Updated(e) => {
1885            ensure_accepted_emitted(
1886                client_order_id,
1887                account_id,
1888                venue_order_id,
1889                identity,
1890                emitter,
1891                state,
1892                ts_init,
1893                ts_init,
1894            );
1895            is_terminal = false;
1896            emitter.send_order_event(OrderEventAny::Updated(e));
1897        }
1898        ParsedOrderEvent::Fill(fill_report) => {
1899            is_terminal = venue_status == OKXOrderStatus::Filled;
1900
1901            if state.check_and_insert_trade(fill_report.trade_id) {
1902                log::debug!(
1903                    "Skipping duplicate fill for {client_order_id}: trade_id={}",
1904                    fill_report.trade_id
1905                );
1906            } else {
1907                emit_venue_order_id_update_if_changed(
1908                    client_order_id,
1909                    account_id,
1910                    venue_order_id,
1911                    identity,
1912                    fill_report.ts_event,
1913                    emitter,
1914                    state,
1915                    order_state_cache,
1916                    ts_init,
1917                );
1918                ensure_accepted_emitted(
1919                    client_order_id,
1920                    account_id,
1921                    venue_order_id,
1922                    identity,
1923                    emitter,
1924                    state,
1925                    ts_init,
1926                    ts_init,
1927                );
1928                state.insert_filled(client_order_id);
1929                state.remove_triggered(&client_order_id);
1930                let filled = fill_report_to_order_filled(
1931                    &fill_report,
1932                    emitter.trader_id(),
1933                    identity,
1934                    instrument.quote_currency(),
1935                );
1936                emitter.send_order_event(OrderEventAny::Filled(filled));
1937            }
1938        }
1939        ParsedOrderEvent::StatusOnly(report) => {
1940            is_terminal = matches!(
1941                report.order_status,
1942                OrderStatus::Filled | OrderStatus::Canceled | OrderStatus::Expired
1943            );
1944            emitter.send_order_status_report(*report);
1945        }
1946        ParsedOrderEvent::Skipped => return,
1947    }
1948
1949    if is_terminal {
1950        state.insert_terminal(client_order_id);
1951        state.remove_order_tracking(client_order_id);
1952        state.remove_accepted(&client_order_id);
1953        order_state_cache.remove(&client_order_id);
1954        // Keep fee_cache and filled_qty_cache entries: replayed terminal
1955        // messages go through the untracked report path and need prior
1956        // cumulative state to avoid re-emitting the full fill quantity
1957    }
1958}
1959
1960/// Synthesizes and emits `OrderAccepted` if one has not yet been emitted for
1961/// this order. Handles fast-filling orders that skip the `Live` state on OKX.
1962#[expect(
1963    clippy::too_many_arguments,
1964    reason = "acceptance reconstruction requires identity, venue state, and event timestamps"
1965)]
1966fn ensure_accepted_emitted(
1967    client_order_id: ClientOrderId,
1968    account_id: AccountId,
1969    venue_order_id: VenueOrderId,
1970    identity: &OrderIdentity,
1971    emitter: &ExecutionEventEmitter,
1972    state: &WsDispatchState,
1973    ts_event: UnixNanos,
1974    ts_init: UnixNanos,
1975) {
1976    if state.contains_accepted(&client_order_id)
1977        || state.contains_terminal(&client_order_id)
1978        || state.contains_filled(&client_order_id)
1979    {
1980        return;
1981    }
1982
1983    state.insert_accepted(client_order_id, venue_order_id);
1984    let accepted = OrderAccepted::new(
1985        emitter.trader_id(),
1986        identity.strategy_id,
1987        identity.instrument_id,
1988        client_order_id,
1989        venue_order_id,
1990        account_id,
1991        UUID4::new(),
1992        ts_event,
1993        ts_init,
1994        false,
1995    );
1996    emitter.send_order_event(OrderEventAny::Accepted(accepted));
1997}
1998
1999#[expect(clippy::too_many_arguments)]
2000fn emit_venue_order_id_update_if_changed(
2001    client_order_id: ClientOrderId,
2002    account_id: AccountId,
2003    venue_order_id: VenueOrderId,
2004    identity: &OrderIdentity,
2005    ts_event: UnixNanos,
2006    emitter: &ExecutionEventEmitter,
2007    state: &WsDispatchState,
2008    order_state_cache: &AHashMap<ClientOrderId, OrderStateSnapshot>,
2009    ts_init: UnixNanos,
2010) {
2011    let Some(accepted_venue_order_id) = state.accepted_venue_order_id(&client_order_id) else {
2012        return;
2013    };
2014
2015    if accepted_venue_order_id == venue_order_id {
2016        return;
2017    }
2018
2019    let Some(snapshot) = order_state_cache.get(&client_order_id) else {
2020        return;
2021    };
2022
2023    state.insert_accepted(client_order_id, venue_order_id);
2024    let updated = OrderUpdated::new(
2025        emitter.trader_id(),
2026        identity.strategy_id,
2027        identity.instrument_id,
2028        client_order_id,
2029        snapshot.quantity,
2030        UUID4::new(),
2031        ts_event,
2032        ts_init,
2033        false,
2034        Some(venue_order_id),
2035        Some(account_id),
2036        snapshot.price,
2037        None,
2038        None,
2039        false,
2040    );
2041    emitter.send_order_event(OrderEventAny::Updated(updated));
2042}
2043
2044/// Converts a [`FillReport`] into an [`OrderFilled`] event using tracked identity.
2045fn fill_report_to_order_filled(
2046    report: &FillReport,
2047    trader_id: TraderId,
2048    identity: &OrderIdentity,
2049    quote_currency: Currency,
2050) -> OrderFilled {
2051    OrderFilled::new(
2052        trader_id,
2053        identity.strategy_id,
2054        report.instrument_id,
2055        report
2056            .client_order_id
2057            .expect("tracked order has client_order_id"),
2058        report.venue_order_id,
2059        report.account_id,
2060        report.trade_id,
2061        identity.order_side,
2062        identity.order_type,
2063        report.last_qty,
2064        report.last_px,
2065        quote_currency,
2066        report.liquidity_side,
2067        UUID4::new(),
2068        report.ts_event,
2069        report.ts_init,
2070        false,
2071        report.venue_position_id,
2072        Some(report.commission),
2073        None,
2074    )
2075}
2076
2077/// Falls back to the report path for a single order message.
2078#[expect(clippy::too_many_arguments)]
2079fn dispatch_order_msg_as_report(
2080    msg: &OKXOrderMsg,
2081    account_id: AccountId,
2082    instruments: &AHashMap<Ustr, InstrumentAny>,
2083    fee_cache: &mut FeeCache,
2084    filled_qty_cache: &mut FilledQtyCache,
2085    emitter: &ExecutionEventEmitter,
2086    state: &WsDispatchState,
2087    ts_init: UnixNanos,
2088) {
2089    match parse_order_msg(
2090        msg,
2091        account_id,
2092        instruments,
2093        fee_cache,
2094        filled_qty_cache,
2095        ts_init,
2096    ) {
2097        Ok(report) => {
2098            dispatch_execution_reports(vec![report], emitter, state);
2099
2100            if let Some(instrument) = instruments.get(&msg.inst_id) {
2101                update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
2102            }
2103        }
2104        Err(e) => log::error!("Failed to parse order message as report: {e}"),
2105    }
2106}
2107
2108#[expect(clippy::too_many_arguments)]
2109fn dispatch_terminal_order_fill_as_report(
2110    msg: &OKXOrderMsg,
2111    client_order_id: ClientOrderId,
2112    account_id: AccountId,
2113    instruments: &AHashMap<Ustr, InstrumentAny>,
2114    fee_cache: &mut FeeCache,
2115    filled_qty_cache: &mut FilledQtyCache,
2116    emitter: &ExecutionEventEmitter,
2117    state: &WsDispatchState,
2118    ts_init: UnixNanos,
2119) {
2120    match parse_order_msg(
2121        msg,
2122        account_id,
2123        instruments,
2124        fee_cache,
2125        filled_qty_cache,
2126        ts_init,
2127    ) {
2128        Ok(ExecutionReport::Fill(mut report)) => {
2129            report.client_order_id = Some(client_order_id);
2130            dispatch_execution_reports(vec![ExecutionReport::Fill(report)], emitter, state);
2131
2132            if let Some(instrument) = instruments.get(&msg.inst_id) {
2133                update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
2134            }
2135        }
2136        Ok(ExecutionReport::Order(_)) => {
2137            log::debug!(
2138                "Suppressing stale regular order status for terminal tracked order {client_order_id}: ord_id={}",
2139                msg.ord_id,
2140            );
2141        }
2142        Err(e) => log::error!("Failed to parse terminal order update: {e}"),
2143    }
2144}
2145
2146fn dispatch_spread_order_msg_as_report(
2147    msg: &OKXSpreadOrder,
2148    account_id: AccountId,
2149    instruments: &AHashMap<Ustr, InstrumentAny>,
2150    filled_qty_cache: &mut FilledQtyCache,
2151    emitter: &ExecutionEventEmitter,
2152    state: &WsDispatchState,
2153    ts_init: UnixNanos,
2154) {
2155    match parse_spread_order_msg(msg, account_id, instruments, filled_qty_cache, ts_init) {
2156        Ok(report) => {
2157            dispatch_execution_reports(vec![report], emitter, state);
2158
2159            if let Some(instrument) = instruments.get(&msg.sprd_id) {
2160                update_spread_fill_cache(msg, instrument, filled_qty_cache);
2161            }
2162        }
2163        Err(e) => log::error!("Failed to parse spread order message as report: {e}"),
2164    }
2165}
2166
2167/// Updates fee, fill, and order state caches from a raw OKX order message.
2168fn update_order_state_cache(
2169    msg: &OKXOrderMsg,
2170    instrument: &InstrumentAny,
2171    client_order_id: ClientOrderId,
2172    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
2173) {
2174    let venue_order_id = VenueOrderId::new(msg.ord_id);
2175    let quantity = parse_quantity(&msg.sz, instrument.size_precision()).unwrap_or_default();
2176    let price = if is_market_price(&msg.px) {
2177        None
2178    } else {
2179        parse_price(&msg.px, instrument.price_precision()).ok()
2180    };
2181
2182    order_state_cache.insert(
2183        client_order_id,
2184        OrderStateSnapshot {
2185            venue_order_id,
2186            quantity,
2187            price,
2188        },
2189    );
2190}
2191
2192fn update_spread_order_state_cache(
2193    msg: &OKXSpreadOrder,
2194    instrument: &InstrumentAny,
2195    client_order_id: ClientOrderId,
2196    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
2197) {
2198    let venue_order_id = VenueOrderId::new(msg.ord_id.as_str());
2199    let quantity = parse_quantity(&msg.sz, instrument.size_precision()).unwrap_or_default();
2200    let price = if is_market_price(&msg.px) {
2201        None
2202    } else {
2203        parse_price(&msg.px, instrument.price_precision()).ok()
2204    };
2205
2206    order_state_cache.insert(
2207        client_order_id,
2208        OrderStateSnapshot {
2209            venue_order_id,
2210            quantity,
2211            price,
2212        },
2213    );
2214}
2215
2216fn update_spread_fill_cache(
2217    msg: &OKXSpreadOrder,
2218    instrument: &InstrumentAny,
2219    filled_qty_cache: &mut FilledQtyCache,
2220) {
2221    if !msg.acc_fill_sz.is_empty()
2222        && msg.acc_fill_sz != "0"
2223        && let Ok(qty) = parse_quantity(&msg.acc_fill_sz, instrument.size_precision())
2224    {
2225        filled_qty_cache.record(msg.ord_id, qty, is_terminal_order_state(msg.state));
2226    }
2227}
2228
2229fn is_spread_post_only_auto_cancel(msg: &OKXSpreadOrder) -> bool {
2230    msg.state == OKXOrderStatus::Canceled && msg.cancel_source == OKX_POST_ONLY_CANCEL_SOURCE
2231}
2232
2233/// Dispatches execution reports with cross-stream deduplication.
2234pub fn dispatch_execution_reports(
2235    reports: Vec<ExecutionReport>,
2236    emitter: &ExecutionEventEmitter,
2237    state: &WsDispatchState,
2238) {
2239    log::debug!("Processing {} execution report(s)", reports.len());
2240
2241    for report in reports {
2242        match report {
2243            ExecutionReport::Order(order_report) => {
2244                if let Some(cid) = order_report.client_order_id {
2245                    match order_report.order_status {
2246                        // Guard form reformats awkwardly across multiple lines
2247                        #[allow(clippy::collapsible_match)]
2248                        OrderStatus::Accepted => {
2249                            if state.contains_terminal(&cid)
2250                                || state.contains_filled(&cid)
2251                                || state.contains_triggered(&cid)
2252                            {
2253                                log::debug!(
2254                                    "Skipping stale OrderStatusReport(Accepted) \
2255                                     for {cid} (order already terminal)"
2256                                );
2257                                continue;
2258                            }
2259
2260                            if !state.contains_accepted(&cid) {
2261                                state.insert_accepted(cid, order_report.venue_order_id);
2262                            }
2263                        }
2264                        OrderStatus::Triggered => {
2265                            if state.contains_filled(&cid) {
2266                                log::debug!(
2267                                    "Skipping stale OrderStatusReport(Triggered) \
2268                                     for {cid} (already filled)"
2269                                );
2270                                continue;
2271                            }
2272                            state.insert_triggered(cid);
2273                        }
2274                        OrderStatus::Filled => {
2275                            state.insert_filled(cid);
2276                            state.insert_terminal(cid);
2277                            state.remove_triggered(&cid);
2278                        }
2279                        OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected => {
2280                            state.insert_terminal(cid);
2281                            state.remove_triggered(&cid);
2282                            state.remove_filled(&cid);
2283                        }
2284                        _ => {}
2285                    }
2286                }
2287                emitter.send_order_status_report(order_report);
2288            }
2289            ExecutionReport::Fill(fill_report) => {
2290                if state.check_and_insert_trade(fill_report.trade_id) {
2291                    log::debug!(
2292                        "Skipping duplicate fill report: trade_id={}",
2293                        fill_report.trade_id
2294                    );
2295                    continue;
2296                }
2297
2298                if let Some(cid) = fill_report.client_order_id {
2299                    state.insert_filled(cid);
2300                    state.remove_triggered(&cid);
2301                }
2302                emitter.send_fill_report(fill_report);
2303            }
2304        }
2305    }
2306}
2307
2308fn emit_send_failed_submit(
2309    failure: &CommandFailure,
2310    state: &WsDispatchState,
2311    emitter: &ExecutionEventEmitter,
2312    clock: &AtomicTime,
2313    client_order_id: ClientOrderId,
2314) {
2315    let (CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason)) = failure else {
2316        return;
2317    };
2318    let Some(ident) = state.order_identity(client_order_id) else {
2319        return;
2320    };
2321
2322    state.remove_order_tracking(client_order_id);
2323    emitter.emit_order_rejected_event(
2324        ident.strategy_id,
2325        ident.instrument_id,
2326        client_order_id,
2327        reason,
2328        clock.get_time_ns(),
2329        false,
2330    );
2331}
2332
2333fn emit_send_failed_modify(
2334    failure: &CommandFailure,
2335    state: &WsDispatchState,
2336    emitter: &ExecutionEventEmitter,
2337    clock: &AtomicTime,
2338    client_order_id: ClientOrderId,
2339) {
2340    let (CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason)) = failure else {
2341        return;
2342    };
2343    let Some(ident) = state.order_identity(client_order_id) else {
2344        return;
2345    };
2346
2347    emitter.emit_order_modify_rejected_event(
2348        ident.strategy_id,
2349        ident.instrument_id,
2350        client_order_id,
2351        None,
2352        reason,
2353        clock.get_time_ns(),
2354    );
2355}
2356
2357fn format_order_response_reason(s_code: &str, s_msg: &str, sub_code: &str) -> String {
2358    match (s_msg.is_empty(), sub_code.is_empty(), s_code.is_empty()) {
2359        (false, true, _) => s_msg.to_string(),
2360        (false, false, _) => format!("{s_msg} (subCode={sub_code})"),
2361        (true, false, false) => format!("sCode={s_code} subCode={sub_code}"),
2362        (true, false, true) => format!("subCode={sub_code}"),
2363        (true, true, false) => format!("sCode={s_code}"),
2364        (true, true, true) => String::new(),
2365    }
2366}
2367
2368#[derive(Debug, Clone)]
2369pub struct AlgoCancelContext {
2370    pub client_order_id: ClientOrderId,
2371    pub instrument_id: InstrumentId,
2372    pub strategy_id: StrategyId,
2373    pub venue_order_id: Option<VenueOrderId>,
2374}
2375
2376// Contexts must correspond 1:1 with the requests that produced
2377// the responses (OKX preserves request order in batch responses).
2378pub fn emit_algo_cancel_rejections(
2379    responses: &[OKXCancelAlgoOrderResponse],
2380    contexts: &[AlgoCancelContext],
2381    emitter: &ExecutionEventEmitter,
2382    clock: &'static AtomicTime,
2383) {
2384    for (i, item) in responses.iter().enumerate() {
2385        let code = item.s_code.as_deref().unwrap_or(OKX_SUCCESS_CODE);
2386        if code == OKX_SUCCESS_CODE {
2387            continue;
2388        }
2389
2390        let msg = item.s_msg.as_deref().unwrap_or("");
2391
2392        if matches!(
2393            classify_okx_venue_code(code, msg),
2394            CommandFailure::Ambiguous(_) | CommandFailure::NotSent(_)
2395        ) {
2396            if let Some(ctx) = contexts.get(i) {
2397                log::warn!(
2398                    "Ambiguous algo cancel response for {}, awaiting reconciliation: \
2399                     algo_id={} sCode={code} sMsg={msg}",
2400                    ctx.client_order_id,
2401                    item.algo_id
2402                );
2403            } else {
2404                log::warn!(
2405                    "Ambiguous algo cancel response without context at index {i}: \
2406                     algo_id={} sCode={code} sMsg={msg}",
2407                    item.algo_id
2408                );
2409            }
2410            continue;
2411        }
2412
2413        if let Some(ctx) = contexts.get(i) {
2414            let ts = clock.get_time_ns();
2415            emitter.emit_order_cancel_rejected_event(
2416                ctx.strategy_id,
2417                ctx.instrument_id,
2418                ctx.client_order_id,
2419                ctx.venue_order_id,
2420                msg,
2421                ts,
2422            );
2423        } else {
2424            log::warn!(
2425                "Algo cancel rejected but no context at index {i}: \
2426                 algo_id={} sCode={code} sMsg={msg}",
2427                item.algo_id
2428            );
2429        }
2430    }
2431}
2432
2433pub fn emit_batch_cancel_failure(
2434    contexts: &[AlgoCancelContext],
2435    error: &str,
2436    _emitter: &ExecutionEventEmitter,
2437    _clock: &'static AtomicTime,
2438) {
2439    for ctx in contexts {
2440        log::warn!(
2441            "Ambiguous algo batch cancel failure for {}, awaiting reconciliation: {error}",
2442            ctx.client_order_id
2443        );
2444    }
2445}
2446
2447#[cfg(test)]
2448mod tests {
2449    use std::time::Duration;
2450
2451    use nautilus_common::messages::{ExecutionEvent, ExecutionReport as CommonExecutionReport};
2452    use nautilus_core::time::get_atomic_clock_realtime;
2453    use nautilus_model::{
2454        enums::{AccountType, OrderSide, OrderType, TimeInForce, TriggerType},
2455        identifiers::Symbol,
2456        instruments::CryptoPerpetual,
2457        types::{Price, Quantity},
2458    };
2459    use rstest::rstest;
2460
2461    use super::*;
2462    use crate::websocket::{error::OKXWsError, messages::OKXWsFrame};
2463
2464    fn load_algo_order_messages(fixture: &str) -> Vec<OKXAlgoOrderMsg> {
2465        let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2466            .join("test_data")
2467            .join(fixture);
2468        let content = std::fs::read_to_string(path).unwrap();
2469        let frame: OKXWsFrame = serde_json::from_str(&content).unwrap();
2470        let OKXWsFrame::Data { data, .. } = frame else {
2471            panic!("Expected algo order data frame");
2472        };
2473        serde_json::from_value(data).unwrap()
2474    }
2475
2476    fn load_regular_order_messages(fixture: &str) -> Vec<OKXOrderMsg> {
2477        let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2478            .join("test_data")
2479            .join(fixture);
2480        let content = std::fs::read_to_string(path).unwrap();
2481        let frame: OKXWsFrame = serde_json::from_str(&content).unwrap();
2482        let OKXWsFrame::Data { data, .. } = frame else {
2483            panic!("Expected regular order data frame");
2484        };
2485        serde_json::from_value(data).unwrap()
2486    }
2487
2488    fn test_algo_context(client_order_id: ClientOrderId) -> OrderContext {
2489        OrderContext {
2490            identity: OrderIdentity {
2491                client_order_id,
2492                strategy_id: StrategyId::from("STRATEGY-001"),
2493                instrument_id: InstrumentId::from("BTC-USDT-SWAP.OKX"),
2494                order_side: OrderSide::Sell,
2495                order_type: OrderType::StopLimit,
2496            },
2497            quantity: Quantity::from("0.01"),
2498            price: Some(Price::from("102900")),
2499            trigger_price: Some(Price::from("95000")),
2500            trigger_type: Some(TriggerType::LastPrice),
2501            time_in_force: TimeInForce::Gtc,
2502            is_post_only: false,
2503            is_reduce_only: true,
2504            is_quote_quantity: false,
2505        }
2506    }
2507
2508    fn test_algo_instruments() -> AtomicMap<Ustr, InstrumentAny> {
2509        let instrument = CryptoPerpetual::builder()
2510            .instrument_id(InstrumentId::from("BTC-USDT-SWAP.OKX"))
2511            .raw_symbol(Symbol::from("BTC-USDT-SWAP"))
2512            .base_currency(Currency::BTC())
2513            .quote_currency(Currency::USDT())
2514            .settlement_currency(Currency::USDT())
2515            .is_inverse(false)
2516            .price_precision(2)
2517            .size_precision(8)
2518            .price_increment(Price::from("0.01"))
2519            .size_increment(Quantity::from("0.00000001"))
2520            .ts_event(UnixNanos::default())
2521            .ts_init(UnixNanos::default())
2522            .build()
2523            .unwrap();
2524        let instruments = AtomicMap::new();
2525        instruments.insert(
2526            Ustr::from("BTC-USDT-SWAP"),
2527            InstrumentAny::CryptoPerpetual(instrument),
2528        );
2529        instruments
2530    }
2531
2532    fn test_execution_emitter() -> (
2533        ExecutionEventEmitter,
2534        tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2535    ) {
2536        let clock = get_atomic_clock_realtime();
2537        let mut emitter = ExecutionEventEmitter::new(
2538            clock,
2539            TraderId::from("TRADER-001"),
2540            AccountId::from("OKX-001"),
2541            AccountType::Margin,
2542            None,
2543        );
2544        let (sender, receiver) = tokio::sync::mpsc::unbounded_channel();
2545        emitter.set_sender(sender);
2546        (emitter, receiver)
2547    }
2548
2549    fn dispatch_test_message(
2550        message: OKXWsMessage,
2551        emitter: &ExecutionEventEmitter,
2552        state: &WsDispatchState,
2553        instruments: &AtomicMap<Ustr, InstrumentAny>,
2554    ) {
2555        dispatch_ws_message(
2556            message,
2557            emitter,
2558            state,
2559            AccountId::from("OKX-001"),
2560            AccountType::Margin,
2561            instruments,
2562            &mut FeeCache::new(),
2563            &mut FilledQtyCache::new(),
2564            &mut AHashMap::new(),
2565            get_atomic_clock_realtime(),
2566        );
2567    }
2568
2569    fn drain_execution_events(
2570        receiver: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2571    ) -> Vec<ExecutionEvent> {
2572        let mut events = Vec::new();
2573        while let Ok(event) = receiver.try_recv() {
2574            events.push(event);
2575        }
2576        events
2577    }
2578
2579    #[rstest]
2580    #[case::margin(AccountType::Margin)]
2581    #[case::cash(AccountType::Cash)]
2582    fn account_updates_report_the_configured_account_type(#[case] account_type: AccountType) {
2583        let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2584            .join("test_data")
2585            .join("http_get_account_balance.json");
2586        let content = std::fs::read_to_string(path).unwrap();
2587        let payload: serde_json::Value = serde_json::from_str(&content).unwrap();
2588        let instruments = AtomicMap::new();
2589        let (emitter, mut receiver) = test_execution_emitter();
2590
2591        dispatch_ws_message(
2592            OKXWsMessage::Account(payload["data"].clone()),
2593            &emitter,
2594            &WsDispatchState::default(),
2595            AccountId::from("OKX-001"),
2596            account_type,
2597            &instruments,
2598            &mut FeeCache::new(),
2599            &mut FilledQtyCache::new(),
2600            &mut AHashMap::new(),
2601            get_atomic_clock_realtime(),
2602        );
2603
2604        let events = drain_execution_events(&mut receiver);
2605        assert_eq!(events.len(), 1);
2606
2607        let ExecutionEvent::Account(state) = &events[0] else {
2608            panic!("expected an account state event");
2609        };
2610
2611        assert_eq!(state.account_type, account_type);
2612        assert_eq!(state.account_id, AccountId::from("OKX-001"));
2613    }
2614
2615    #[rstest]
2616    fn tracked_algo_live_and_pause_emit_one_parent_acceptance() {
2617        let mut messages = load_algo_order_messages("ws_orders_algo.json");
2618        let live = messages.remove(0);
2619        let mut pause = live.clone();
2620        pause.state = OKXAlgoOrderStatus::Pause;
2621        pause.u_time += 1;
2622        let client_order_id = ClientOrderId::new(live.algo_cl_ord_id.as_str());
2623        let state = WsDispatchState::default();
2624        state.track_order_context(test_algo_context(client_order_id));
2625        let instruments = test_algo_instruments();
2626        let (emitter, mut receiver) = test_execution_emitter();
2627
2628        dispatch_test_message(
2629            OKXWsMessage::AlgoOrders(vec![live, pause]),
2630            &emitter,
2631            &state,
2632            &instruments,
2633        );
2634
2635        let events = drain_execution_events(&mut receiver);
2636        assert_eq!(events.len(), 1);
2637        match &events[0] {
2638            ExecutionEvent::Order(OrderEventAny::Accepted(accepted)) => {
2639                assert_eq!(accepted.client_order_id, client_order_id);
2640                assert_eq!(
2641                    accepted.venue_order_id,
2642                    VenueOrderId::new("706620792746729472")
2643                );
2644            }
2645            other => panic!("Expected tracked algo acceptance, was {other:?}"),
2646        }
2647    }
2648
2649    #[rstest]
2650    fn algo_update_routes_are_explicit_for_external_tracked_and_suppressed() {
2651        let mut message = load_algo_order_messages("ws_orders_algo.json").remove(0);
2652        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
2653        let state = WsDispatchState::default();
2654
2655        assert_eq!(
2656            route_algo_order_message(&message, &state),
2657            ExecutionUpdateRoute::External
2658        );
2659
2660        let context = test_algo_context(client_order_id);
2661        state.track_order_context(context);
2662        assert_eq!(
2663            route_algo_order_message(&message, &state),
2664            ExecutionUpdateRoute::Tracked(client_order_id, context)
2665        );
2666
2667        message.state = OKXAlgoOrderStatus::Unknown;
2668        assert_eq!(
2669            route_algo_order_message(&message, &state),
2670            ExecutionUpdateRoute::Suppressed
2671        );
2672    }
2673
2674    #[rstest]
2675    fn rest_parent_binding_routes_linked_child_before_parent_stream_update() {
2676        let mut child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2677        let client_order_id = ClientOrderId::new("STOP003BTCUSDT20250120");
2678        let parent_venue_order_id = VenueOrderId::new("706620792746729474");
2679        let child_venue_order_id = VenueOrderId::new(child.ord_id);
2680        child.algo_cl_ord_id = None;
2681        let state = WsDispatchState::default();
2682        state.track_order_context(test_algo_context(client_order_id));
2683        state.bind_algo_parent(client_order_id, parent_venue_order_id);
2684        let instruments = test_algo_instruments();
2685        let (emitter, mut receiver) = test_execution_emitter();
2686
2687        dispatch_test_message(
2688            OKXWsMessage::Orders(vec![child.clone()]),
2689            &emitter,
2690            &state,
2691            &instruments,
2692        );
2693
2694        let events = drain_execution_events(&mut receiver);
2695        assert_eq!(events.len(), 4);
2696        assert!(matches!(
2697            &events[0],
2698            ExecutionEvent::Order(OrderEventAny::Accepted(accepted))
2699                if accepted.client_order_id == client_order_id
2700                    && accepted.venue_order_id == parent_venue_order_id
2701        ));
2702        assert!(matches!(
2703            &events[1],
2704            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2705                if updated.client_order_id == client_order_id
2706                    && updated.venue_order_id == Some(child_venue_order_id)
2707        ));
2708        assert!(matches!(
2709            &events[2],
2710            ExecutionEvent::Order(OrderEventAny::Triggered(triggered))
2711                if triggered.client_order_id == client_order_id
2712                    && triggered.venue_order_id == Some(child_venue_order_id)
2713        ));
2714        assert!(matches!(
2715            &events[3],
2716            ExecutionEvent::Order(OrderEventAny::Filled(filled))
2717                if filled.client_order_id == client_order_id
2718                    && filled.venue_order_id == child_venue_order_id
2719        ));
2720
2721        dispatch_test_message(
2722            OKXWsMessage::Orders(vec![child]),
2723            &emitter,
2724            &state,
2725            &instruments,
2726        );
2727        assert!(drain_execution_events(&mut receiver).is_empty());
2728    }
2729
2730    #[rstest]
2731    #[tokio::test]
2732    async fn pre_binding_child_is_held_until_parent_binding() {
2733        let child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2734        let parent_venue_order_id = VenueOrderId::new(
2735            child
2736                .linked_algo_ord
2737                .as_ref()
2738                .expect("Expected linked algo order")
2739                .algo_id
2740                .as_str(),
2741        );
2742        let client_order_id = ClientOrderId::new("STOP003BTCUSDT20250120");
2743        let child_venue_order_id = VenueOrderId::new(child.ord_id);
2744        let state = WsDispatchState::default();
2745        state.track_order_context(test_algo_context(client_order_id));
2746        let instruments = test_algo_instruments();
2747        let (emitter, mut receiver) = test_execution_emitter();
2748
2749        dispatch_test_message(
2750            OKXWsMessage::Orders(vec![child.clone()]),
2751            &emitter,
2752            &state,
2753            &instruments,
2754        );
2755        assert!(drain_execution_events(&mut receiver).is_empty());
2756
2757        state.bind_algo_parent(client_order_id, parent_venue_order_id);
2758        tokio::time::timeout(
2759            Duration::from_millis(100),
2760            state.wait_for_linked_child_route(),
2761        )
2762        .await
2763        .expect("Expected parent binding to wake the private dispatch loop");
2764        dispatch_test_message(
2765            OKXWsMessage::Orders(Vec::new()),
2766            &emitter,
2767            &state,
2768            &instruments,
2769        );
2770
2771        let events = drain_execution_events(&mut receiver);
2772        assert_eq!(events.len(), 4);
2773        assert!(matches!(
2774            &events[0],
2775            ExecutionEvent::Order(OrderEventAny::Accepted(accepted))
2776                if accepted.client_order_id == client_order_id
2777                    && accepted.venue_order_id == parent_venue_order_id
2778        ));
2779        assert!(matches!(
2780            &events[1],
2781            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2782                if updated.client_order_id == client_order_id
2783                    && updated.venue_order_id == Some(child_venue_order_id)
2784        ));
2785        assert!(matches!(
2786            &events[2],
2787            ExecutionEvent::Order(OrderEventAny::Triggered(triggered))
2788                if triggered.client_order_id == client_order_id
2789                    && triggered.venue_order_id == Some(child_venue_order_id)
2790        ));
2791        assert!(matches!(
2792            &events[3],
2793            ExecutionEvent::Order(OrderEventAny::Filled(filled))
2794                if filled.client_order_id == client_order_id
2795                    && filled.venue_order_id == child_venue_order_id
2796        ));
2797
2798        dispatch_test_message(OKXWsMessage::Reconnected, &emitter, &state, &instruments);
2799        dispatch_test_message(
2800            OKXWsMessage::Orders(vec![child]),
2801            &emitter,
2802            &state,
2803            &instruments,
2804        );
2805        assert!(drain_execution_events(&mut receiver).is_empty());
2806    }
2807
2808    #[rstest]
2809    fn held_child_becomes_external_when_candidate_binds_another_parent() {
2810        let child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2811        let client_order_id = ClientOrderId::new("STOP003BTCUSDT20250120");
2812        let child_venue_order_id = VenueOrderId::new(child.ord_id);
2813        let state = WsDispatchState::default();
2814        state.track_order_context(test_algo_context(client_order_id));
2815        let instruments = test_algo_instruments();
2816        let (emitter, mut receiver) = test_execution_emitter();
2817
2818        dispatch_test_message(
2819            OKXWsMessage::Orders(vec![child]),
2820            &emitter,
2821            &state,
2822            &instruments,
2823        );
2824        assert!(drain_execution_events(&mut receiver).is_empty());
2825
2826        let authoritative_parent = VenueOrderId::new("706620792746729475");
2827        state.bind_algo_parent(client_order_id, authoritative_parent);
2828        dispatch_test_message(
2829            OKXWsMessage::Orders(Vec::new()),
2830            &emitter,
2831            &state,
2832            &instruments,
2833        );
2834
2835        let events = drain_execution_events(&mut receiver);
2836        assert_eq!(events.len(), 1);
2837        assert!(matches!(
2838            &events[0],
2839            ExecutionEvent::Report(CommonExecutionReport::Fill(report))
2840                if report.client_order_id
2841                    == Some(ClientOrderId::new("706620792746729474_0"))
2842                    && report.venue_order_id == child_venue_order_id
2843        ));
2844        assert_eq!(
2845            state.order_venue_binding(client_order_id),
2846            Some((authoritative_parent, false))
2847        );
2848    }
2849
2850    #[rstest]
2851    fn active_algo_bindings_are_not_evicted_with_replay_state() {
2852        let state = WsDispatchState::default();
2853        let first_client_order_id = ClientOrderId::new("STOP-ACTIVE-00000");
2854        let first_venue_order_id = VenueOrderId::new("700000000000000000");
2855
2856        for index in 0..=DEDUP_CAPACITY {
2857            let client_order_id = ClientOrderId::new(format!("STOP-ACTIVE-{index:05}").as_str());
2858            let venue_order_id = VenueOrderId::new(
2859                format!("{}", 700_000_000_000_000_000_u64 + index as u64).as_str(),
2860            );
2861            state.bind_algo_parent(client_order_id, venue_order_id);
2862        }
2863
2864        assert_eq!(
2865            state.order_venue_binding(first_client_order_id),
2866            Some((first_venue_order_id, false))
2867        );
2868    }
2869
2870    #[rstest]
2871    fn ambiguous_submit_failure_preserves_only_bound_context() {
2872        let client_order_id = ClientOrderId::new("STOP-AMBIGUOUS-001");
2873        let failure = CommandFailure::Ambiguous("request timed out".to_string());
2874        let unbound_state = WsDispatchState::default();
2875        unbound_state.track_order_context(test_algo_context(client_order_id));
2876
2877        unbound_state.resolve_algo_submit_failure(client_order_id, &failure);
2878        assert_eq!(unbound_state.order_identity(client_order_id), None);
2879
2880        let bound_state = WsDispatchState::default();
2881        let context = test_algo_context(client_order_id);
2882        let parent_venue_order_id = VenueOrderId::new("706620792746729476");
2883        bound_state.track_order_context(context);
2884        bound_state.bind_algo_parent(client_order_id, parent_venue_order_id);
2885
2886        bound_state.resolve_algo_submit_failure(client_order_id, &failure);
2887        assert_eq!(
2888            bound_state.order_identity(client_order_id),
2889            Some(context.identity)
2890        );
2891        assert_eq!(
2892            bound_state.order_venue_binding(client_order_id),
2893            Some((parent_venue_order_id, false))
2894        );
2895    }
2896
2897    #[rstest]
2898    #[case::effective(0, false)]
2899    #[case::effective_order_id_list(0, true)]
2900    #[case::partially_effective(1, false)]
2901    #[case::order_placed(2, false)]
2902    fn tracked_algo_trigger_states_bind_child_before_trigger(
2903        #[case] index: usize,
2904        #[case] use_order_id_list: bool,
2905    ) {
2906        let mut messages = if index == 2 {
2907            load_algo_order_messages("ws_orders_algo.json")
2908        } else {
2909            load_algo_order_messages("ws_orders_algo_states.json")
2910        };
2911        let mut message = if index == 2 {
2912            messages.remove(2)
2913        } else {
2914            messages.remove(index)
2915        };
2916        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
2917        let child_venue_order_id = VenueOrderId::new(message.ord_id.as_str());
2918        if use_order_id_list {
2919            message.ord_id_list = vec![message.ord_id.clone()];
2920            message.ord_id.clear();
2921        }
2922
2923        let state = WsDispatchState::default();
2924        state.track_order_context(test_algo_context(client_order_id));
2925        let instruments = test_algo_instruments();
2926        let (emitter, mut receiver) = test_execution_emitter();
2927
2928        dispatch_test_message(
2929            OKXWsMessage::AlgoOrders(vec![message]),
2930            &emitter,
2931            &state,
2932            &instruments,
2933        );
2934
2935        let events = drain_execution_events(&mut receiver);
2936        assert_eq!(events.len(), 3);
2937        assert!(matches!(
2938            &events[0],
2939            ExecutionEvent::Order(OrderEventAny::Accepted(_))
2940        ));
2941        assert!(matches!(
2942            &events[1],
2943            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2944                if updated.venue_order_id == Some(child_venue_order_id)
2945        ));
2946        assert!(matches!(
2947            &events[2],
2948            ExecutionEvent::Order(OrderEventAny::Triggered(triggered))
2949                if triggered.venue_order_id == Some(child_venue_order_id)
2950        ));
2951        assert_eq!(
2952            state.order_venue_binding(client_order_id),
2953            Some((child_venue_order_id, true))
2954        );
2955    }
2956
2957    #[rstest]
2958    fn tracked_algo_transition_uses_venue_actual_quantity() {
2959        let mut message = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
2960        message.sz.clear();
2961        message.close_fraction = "1".to_string();
2962        message.actual_sz = "0.025".to_string();
2963        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
2964        let state = WsDispatchState::default();
2965        state.track_order_context(test_algo_context(client_order_id));
2966        let instruments = test_algo_instruments();
2967        let (emitter, mut receiver) = test_execution_emitter();
2968
2969        dispatch_test_message(
2970            OKXWsMessage::AlgoOrders(vec![message]),
2971            &emitter,
2972            &state,
2973            &instruments,
2974        );
2975
2976        let events = drain_execution_events(&mut receiver);
2977        assert_eq!(events.len(), 3);
2978        assert!(matches!(
2979            &events[1],
2980            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2981                if updated.quantity == Quantity::from("0.025")
2982                    && updated.price.is_none()
2983        ));
2984        assert_eq!(
2985            state
2986                .order_contexts
2987                .get(&client_order_id)
2988                .map(|context| context.quantity),
2989            Some(Quantity::from("0.025"))
2990        );
2991    }
2992
2993    #[rstest]
2994    fn regular_child_refreshes_parent_transition_terms() {
2995        let mut parent = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
2996        parent.ord_px = "94950".to_string();
2997        let parent_replay = parent.clone();
2998        let client_order_id = ClientOrderId::new(parent.algo_cl_ord_id.as_str());
2999        let parent_venue_order_id = parent.algo_id.clone();
3000        let child_venue_order_id = parent.ord_id.clone();
3001        let state = WsDispatchState::default();
3002        state.track_order_context(test_algo_context(client_order_id));
3003        let instruments = test_algo_instruments();
3004        let (emitter, mut receiver) = test_execution_emitter();
3005
3006        dispatch_test_message(
3007            OKXWsMessage::AlgoOrders(vec![parent]),
3008            &emitter,
3009            &state,
3010            &instruments,
3011        );
3012        assert_eq!(drain_execution_events(&mut receiver).len(), 3);
3013
3014        let mut child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
3015        child.algo_id = Some(parent_venue_order_id.clone());
3016        child.algo_cl_ord_id = None;
3017        child.linked_algo_ord = Some(crate::websocket::messages::OKXLinkedAlgoOrd {
3018            algo_id: parent_venue_order_id,
3019        });
3020        child.ord_id = Ustr::from(child_venue_order_id.as_str());
3021        child.ord_type = OKXOrderType::Limit;
3022        child.state = OKXOrderStatus::Live;
3023        child.sz = "0.025".to_string();
3024        child.px = "94950".to_string();
3025        child.acc_fill_sz = Some("0".to_string());
3026        child.fill_sz.clear();
3027        child.fill_px.clear();
3028        child.trade_id.clear();
3029
3030        dispatch_test_message(
3031            OKXWsMessage::Orders(vec![child]),
3032            &emitter,
3033            &state,
3034            &instruments,
3035        );
3036
3037        let events = drain_execution_events(&mut receiver);
3038        assert_eq!(events.len(), 1);
3039        assert!(matches!(
3040            &events[0],
3041            ExecutionEvent::Order(OrderEventAny::Updated(updated))
3042                if updated.quantity == Quantity::from("0.025")
3043                    && updated.price == Some(Price::from("94950"))
3044        ));
3045        assert_eq!(
3046            state
3047                .order_contexts
3048                .get(&client_order_id)
3049                .map(|context| (context.quantity, context.price)),
3050            Some((Quantity::from("0.025"), Some(Price::from("94950"))))
3051        );
3052
3053        dispatch_test_message(
3054            OKXWsMessage::AlgoOrders(vec![parent_replay]),
3055            &emitter,
3056            &state,
3057            &instruments,
3058        );
3059        assert!(drain_execution_events(&mut receiver).is_empty());
3060        assert_eq!(
3061            state
3062                .order_contexts
3063                .get(&client_order_id)
3064                .map(|context| context.quantity),
3065            Some(Quantity::from("0.025"))
3066        );
3067    }
3068
3069    #[rstest]
3070    fn tracked_algo_holds_trigger_until_linked_child_arrives() {
3071        let mut parent = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
3072        parent.algo_id = "706620792746729474".to_string();
3073        parent.algo_cl_ord_id = "STOP003BTCUSDT20250120".to_string();
3074        parent.ord_id.clear();
3075        parent.ord_id_list.clear();
3076        let client_order_id = ClientOrderId::new(parent.algo_cl_ord_id.as_str());
3077        let state = WsDispatchState::default();
3078        state.track_order_context(test_algo_context(client_order_id));
3079        let instruments = test_algo_instruments();
3080        let (emitter, mut receiver) = test_execution_emitter();
3081
3082        dispatch_test_message(
3083            OKXWsMessage::AlgoOrders(vec![parent]),
3084            &emitter,
3085            &state,
3086            &instruments,
3087        );
3088
3089        let parent_events = drain_execution_events(&mut receiver);
3090        assert_eq!(parent_events.len(), 1);
3091        assert!(matches!(
3092            &parent_events[0],
3093            ExecutionEvent::Order(OrderEventAny::Accepted(_))
3094        ));
3095        assert!(!state.contains_triggered(&client_order_id));
3096
3097        let child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
3098        assert!(child.algo_cl_ord_id.is_none());
3099        assert_eq!(child.cl_ord_id, "706620792746729474_0");
3100        assert_eq!(
3101            child
3102                .linked_algo_ord
3103                .as_ref()
3104                .map(|linked| linked.algo_id.as_str()),
3105            Some("706620792746729474")
3106        );
3107        dispatch_test_message(
3108            OKXWsMessage::Orders(vec![child]),
3109            &emitter,
3110            &state,
3111            &instruments,
3112        );
3113
3114        let child_events = drain_execution_events(&mut receiver);
3115        assert_eq!(child_events.len(), 3);
3116        assert!(matches!(
3117            &child_events[0],
3118            ExecutionEvent::Order(OrderEventAny::Updated(_))
3119        ));
3120        assert!(matches!(
3121            &child_events[1],
3122            ExecutionEvent::Order(OrderEventAny::Triggered(_))
3123        ));
3124
3125        match &child_events[2] {
3126            ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
3127                assert_eq!(filled.client_order_id, client_order_id);
3128                assert_eq!(
3129                    filled.venue_order_id,
3130                    VenueOrderId::new("706620792746729999")
3131                );
3132                assert_eq!(filled.trade_id, TradeId::new("1518905530"));
3133            }
3134            other => panic!("Expected tracked child fill, was {other:?}"),
3135        }
3136
3137        assert!(state.contains_terminal(&client_order_id));
3138    }
3139
3140    #[rstest]
3141    fn tracked_algo_child_cancellation_routes_late_fill_report() {
3142        let parent = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
3143        let client_order_id = ClientOrderId::new(parent.algo_cl_ord_id.as_str());
3144        let parent_venue_order_id = parent.algo_id.clone();
3145        let child_venue_order_id = parent.ord_id.clone();
3146        let state = WsDispatchState::default();
3147        state.track_order_context(test_algo_context(client_order_id));
3148        let instruments = test_algo_instruments();
3149        let (emitter, mut receiver) = test_execution_emitter();
3150
3151        dispatch_test_message(
3152            OKXWsMessage::AlgoOrders(vec![parent]),
3153            &emitter,
3154            &state,
3155            &instruments,
3156        );
3157        assert_eq!(drain_execution_events(&mut receiver).len(), 3);
3158
3159        let mut child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
3160        child.algo_id = Some(parent_venue_order_id.clone());
3161        child.linked_algo_ord = Some(crate::websocket::messages::OKXLinkedAlgoOrd {
3162            algo_id: parent_venue_order_id,
3163        });
3164        child.ord_id = Ustr::from(child_venue_order_id.as_str());
3165        let mut canceled = child.clone();
3166        canceled.state = OKXOrderStatus::Canceled;
3167        canceled.acc_fill_sz = Some("0".to_string());
3168        canceled.fill_sz.clear();
3169        canceled.fill_px.clear();
3170        canceled.trade_id.clear();
3171        dispatch_test_message(
3172            OKXWsMessage::Orders(vec![canceled]),
3173            &emitter,
3174            &state,
3175            &instruments,
3176        );
3177
3178        let events = drain_execution_events(&mut receiver);
3179        assert_eq!(events.len(), 1);
3180        assert!(matches!(
3181            &events[0],
3182            ExecutionEvent::Order(OrderEventAny::Canceled(canceled))
3183                if canceled.client_order_id == client_order_id
3184                    && canceled.venue_order_id
3185                        == Some(VenueOrderId::new(child_venue_order_id.as_str()))
3186        ));
3187        assert!(state.contains_terminal(&client_order_id));
3188
3189        dispatch_test_message(
3190            OKXWsMessage::Orders(vec![child.clone()]),
3191            &emitter,
3192            &state,
3193            &instruments,
3194        );
3195        let events = drain_execution_events(&mut receiver);
3196        assert_eq!(events.len(), 1);
3197        assert!(matches!(
3198            &events[0],
3199            ExecutionEvent::Report(CommonExecutionReport::Fill(report))
3200                if report.client_order_id == Some(client_order_id)
3201                    && report.venue_order_id
3202                        == VenueOrderId::new(child_venue_order_id.as_str())
3203                    && report.trade_id == TradeId::new("1518905530")
3204        ));
3205
3206        dispatch_test_message(
3207            OKXWsMessage::Orders(vec![child]),
3208            &emitter,
3209            &state,
3210            &instruments,
3211        );
3212        assert!(drain_execution_events(&mut receiver).is_empty());
3213    }
3214
3215    #[rstest]
3216    fn reconnect_replay_and_stale_parent_acceptance_are_suppressed() {
3217        let message = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
3218        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3219        let state = WsDispatchState::default();
3220        state.track_order_context(test_algo_context(client_order_id));
3221        let instruments = test_algo_instruments();
3222        let (emitter, mut receiver) = test_execution_emitter();
3223
3224        dispatch_test_message(
3225            OKXWsMessage::AlgoOrders(vec![message.clone()]),
3226            &emitter,
3227            &state,
3228            &instruments,
3229        );
3230        assert_eq!(drain_execution_events(&mut receiver).len(), 3);
3231
3232        let mut stale_live = message.clone();
3233        stale_live.state = OKXAlgoOrderStatus::Live;
3234        stale_live.ord_id.clear();
3235        dispatch_test_message(OKXWsMessage::Reconnected, &emitter, &state, &instruments);
3236        dispatch_test_message(
3237            OKXWsMessage::AlgoOrders(vec![message, stale_live]),
3238            &emitter,
3239            &state,
3240            &instruments,
3241        );
3242
3243        assert!(drain_execution_events(&mut receiver).is_empty());
3244        assert_eq!(
3245            state
3246                .order_venue_binding(client_order_id)
3247                .map(|(venue_order_id, _)| venue_order_id),
3248            Some(VenueOrderId::new("706620792746730010"))
3249        );
3250    }
3251
3252    #[rstest]
3253    #[case::filled("ws_orders_algo.json", 4, 3)]
3254    #[case::partially_failed("ws_orders_algo_states.json", 4, 1)]
3255    fn tracked_algo_aggregate_state_never_enters_reconciliation(
3256        #[case] fixture: &str,
3257        #[case] index: usize,
3258        #[case] expected_order_events: usize,
3259    ) {
3260        let message = load_algo_order_messages(fixture).remove(index);
3261        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3262        let state = WsDispatchState::default();
3263        state.track_order_context(test_algo_context(client_order_id));
3264        let instruments = test_algo_instruments();
3265        let (emitter, mut receiver) = test_execution_emitter();
3266
3267        dispatch_test_message(
3268            OKXWsMessage::AlgoOrders(vec![message]),
3269            &emitter,
3270            &state,
3271            &instruments,
3272        );
3273
3274        let events = drain_execution_events(&mut receiver);
3275        assert_eq!(events.len(), expected_order_events);
3276        assert!(
3277            events
3278                .iter()
3279                .all(|event| matches!(event, ExecutionEvent::Order(_)))
3280        );
3281        assert!(state.order_contexts.contains_key(&client_order_id));
3282        assert!(!state.contains_terminal(&client_order_id));
3283    }
3284
3285    #[rstest]
3286    fn tracked_algo_cancel_is_terminal_and_stale_live_is_suppressed() {
3287        let mut message = load_algo_order_messages("ws_orders_algo.json").remove(3);
3288        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3289        let state = WsDispatchState::default();
3290        state.track_order_context(test_algo_context(client_order_id));
3291        let instruments = test_algo_instruments();
3292        let (emitter, mut receiver) = test_execution_emitter();
3293
3294        dispatch_test_message(
3295            OKXWsMessage::AlgoOrders(vec![message.clone()]),
3296            &emitter,
3297            &state,
3298            &instruments,
3299        );
3300        message.state = OKXAlgoOrderStatus::Live;
3301        message.algo_cl_ord_id.clear();
3302        message.cl_ord_id.clear();
3303        dispatch_test_message(
3304            OKXWsMessage::AlgoOrders(vec![message]),
3305            &emitter,
3306            &state,
3307            &instruments,
3308        );
3309
3310        let events = drain_execution_events(&mut receiver);
3311        assert_eq!(events.len(), 2);
3312        assert!(matches!(
3313            &events[0],
3314            ExecutionEvent::Order(OrderEventAny::Accepted(_))
3315        ));
3316        assert!(matches!(
3317            &events[1],
3318            ExecutionEvent::Order(OrderEventAny::Canceled(canceled))
3319                if canceled.client_order_id == client_order_id
3320        ));
3321        assert!(state.contains_terminal(&client_order_id));
3322        assert!(!state.order_contexts.contains_key(&client_order_id));
3323        assert_eq!(state.order_venue_binding(client_order_id), None);
3324    }
3325
3326    #[rstest]
3327    fn tracked_algo_failure_uses_typed_rejection() {
3328        let mut message = load_algo_order_messages("ws_orders_algo_states.json").remove(3);
3329        message.fail_code = "51008".to_string();
3330        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3331        let state = WsDispatchState::default();
3332        state.track_order_context(test_algo_context(client_order_id));
3333        let instruments = test_algo_instruments();
3334        let (emitter, mut receiver) = test_execution_emitter();
3335
3336        dispatch_test_message(
3337            OKXWsMessage::AlgoOrders(vec![message]),
3338            &emitter,
3339            &state,
3340            &instruments,
3341        );
3342
3343        let events = drain_execution_events(&mut receiver);
3344        assert_eq!(events.len(), 1);
3345        assert!(matches!(
3346            &events[0],
3347            ExecutionEvent::Order(OrderEventAny::Rejected(rejected))
3348                if rejected.client_order_id == client_order_id
3349                    && rejected.reason == Ustr::from("51008")
3350        ));
3351        assert!(state.contains_terminal(&client_order_id));
3352    }
3353
3354    #[rstest]
3355    #[case("51000", "Rejected", "", "Rejected")]
3356    #[case("51000", "Rejected", "51004", "Rejected (subCode=51004)")]
3357    #[case("51000", "", "51004", "sCode=51000 subCode=51004")]
3358    #[case("51000", "", "", "sCode=51000")]
3359    #[case("", "", "51004", "subCode=51004")]
3360    #[case("", "", "", "")]
3361    fn test_format_order_response_reason(
3362        #[case] s_code: &str,
3363        #[case] s_msg: &str,
3364        #[case] sub_code: &str,
3365        #[case] expected: &str,
3366    ) {
3367        assert_eq!(
3368            format_order_response_reason(s_code, s_msg, sub_code),
3369            expected
3370        );
3371    }
3372
3373    #[rstest]
3374    #[case("order", false, false)]
3375    #[case("batch-orders", false, false)]
3376    #[case("batch-orders", false, true)]
3377    #[case("amend-order", true, false)]
3378    #[case("batch-amend-orders", true, false)]
3379    #[case("batch-amend-orders", true, true)]
3380    fn rpi_minimum_notional_rejection_preserves_order_lifecycle(
3381        #[case] case: &str,
3382        #[case] amend: bool,
3383        #[case] reverse: bool,
3384    ) {
3385        let fixtures: serde_json::Value =
3386            serde_json::from_str(include_str!("../../test_data/rpi_minimum_notional.json"))
3387                .unwrap();
3388        let mut response = fixtures[case]["response"].clone();
3389        if reverse {
3390            response["data"].as_array_mut().unwrap().reverse();
3391        }
3392        let frame: OKXWsFrame = serde_json::from_value(response.clone()).unwrap();
3393        let OKXWsFrame::OrderResponse {
3394            id,
3395            op,
3396            code,
3397            msg,
3398            data,
3399        } = frame
3400        else {
3401            panic!("Expected order response");
3402        };
3403        let state = WsDispatchState::default();
3404        let instrument_id = InstrumentId::from("BTC-USDT.OKX");
3405        let strategy_id = StrategyId::from("RPI-003");
3406        let pending = if amend {
3407            &state.pending_amends
3408        } else {
3409            &state.pending_orders
3410        };
3411        let unrelated = "ORPI003";
3412        pending.insert(
3413            unrelated.to_string(),
3414            PendingOrderInfo {
3415                trader_id: TraderId::from("TRADER-001"),
3416                strategy_id,
3417                instrument_id,
3418            },
3419        );
3420        let mut originals = Vec::new();
3421        for item in &data {
3422            let client_order_id = ClientOrderId::new(item["clOrdId"].as_str().unwrap());
3423            let context = OrderContext {
3424                identity: OrderIdentity {
3425                    client_order_id,
3426                    strategy_id,
3427                    instrument_id,
3428                    order_side: OrderSide::Sell,
3429                    order_type: OrderType::Limit,
3430                },
3431                quantity: Quantity::from("0.2"),
3432                price: Some(Price::from("65123")),
3433                trigger_price: None,
3434                trigger_type: None,
3435                time_in_force: TimeInForce::Gtc,
3436                is_post_only: true,
3437                is_reduce_only: false,
3438                is_quote_quantity: false,
3439            };
3440            state.track_order_context(context);
3441            state
3442                .order_identities
3443                .insert(client_order_id, context.identity);
3444            pending.insert(
3445                client_order_id.to_string(),
3446                PendingOrderInfo {
3447                    trader_id: TraderId::from("TRADER-001"),
3448                    strategy_id,
3449                    instrument_id,
3450                },
3451            );
3452            originals.push(context);
3453        }
3454        let (emitter, mut receiver) = test_execution_emitter();
3455
3456        dispatch_test_message(
3457            OKXWsMessage::OrderResponse {
3458                id,
3459                op,
3460                code,
3461                msg,
3462                data,
3463            },
3464            &emitter,
3465            &state,
3466            &AtomicMap::new(),
3467        );
3468
3469        let events = drain_execution_events(&mut receiver);
3470        assert_eq!(events.len(), 1);
3471        let rejected_id = ClientOrderId::from("ORPI002");
3472        let reason = response["data"]
3473            .as_array()
3474            .unwrap()
3475            .iter()
3476            .find(|item| item["sCode"] == "54051")
3477            .unwrap()["sMsg"]
3478            .as_str()
3479            .unwrap();
3480
3481        match &events[0] {
3482            ExecutionEvent::Order(OrderEventAny::Rejected(event)) if !amend => {
3483                assert_eq!(event.trader_id, TraderId::from("TRADER-001"));
3484                assert_eq!(event.strategy_id, strategy_id);
3485                assert_eq!(event.instrument_id, instrument_id);
3486                assert_eq!(event.client_order_id, rejected_id);
3487                assert_eq!(event.account_id, AccountId::from("OKX-001"));
3488                assert_eq!(event.reason.as_str(), reason);
3489                assert!(!event.reconciliation);
3490                assert!(!event.due_post_only);
3491            }
3492            ExecutionEvent::Order(OrderEventAny::ModifyRejected(event)) if amend => {
3493                assert_eq!(event.trader_id, TraderId::from("TRADER-001"));
3494                assert_eq!(event.strategy_id, strategy_id);
3495                assert_eq!(event.instrument_id, instrument_id);
3496                assert_eq!(event.client_order_id, rejected_id);
3497                assert_eq!(event.account_id, Some(AccountId::from("OKX-001")));
3498                assert_eq!(
3499                    event.venue_order_id,
3500                    Some(VenueOrderId::from("2500000000000000002"))
3501                );
3502                assert_eq!(event.reason.as_str(), reason);
3503                assert!(!event.reconciliation);
3504            }
3505            event => panic!("Unexpected event: {event:?}"),
3506        }
3507        assert_eq!(pending.len(), 1);
3508        assert!(pending.contains_key(unrelated));
3509
3510        for original in originals {
3511            let client_order_id = original.identity.client_order_id;
3512            if amend || client_order_id != rejected_id {
3513                assert_eq!(
3514                    *state.order_identities.get(&client_order_id).unwrap(),
3515                    original.identity
3516                );
3517                assert_eq!(
3518                    *state.order_contexts.get(&client_order_id).unwrap(),
3519                    original
3520                );
3521                assert!(!state.terminal_orders.contains(&client_order_id));
3522            } else {
3523                assert!(!state.order_identities.contains_key(&client_order_id));
3524                assert!(!state.order_contexts.contains_key(&client_order_id));
3525            }
3526        }
3527    }
3528
3529    #[rstest]
3530    #[case::ambiguous(OKXWsError::SendFailed("connection reset".to_string()), true)]
3531    #[case::not_sent(OKXWsError::NoActiveClient, false)]
3532    fn send_failure_preserves_only_ambiguous_pending_orders(
3533        #[case] error: OKXWsError,
3534        #[case] expected_pending: bool,
3535    ) {
3536        let client_order_ids = [
3537            ClientOrderId::from("O-batch-pending-1"),
3538            ClientOrderId::from("O-batch-pending-2"),
3539        ];
3540        let state = WsDispatchState::default();
3541        let instrument_id = InstrumentId::from("ETH-USDT-SWAP.OKX");
3542        let strategy_id = StrategyId::from("STRATEGY-001");
3543
3544        for client_order_id in client_order_ids {
3545            state.pending_orders.insert(
3546                client_order_id.to_string(),
3547                PendingOrderInfo {
3548                    trader_id: TraderId::from("TRADER-001"),
3549                    strategy_id,
3550                    instrument_id,
3551                },
3552            );
3553            state.order_identities.insert(
3554                client_order_id,
3555                OrderIdentity {
3556                    client_order_id,
3557                    instrument_id,
3558                    strategy_id,
3559                    order_side: OrderSide::Buy,
3560                    order_type: OrderType::Limit,
3561                },
3562            );
3563        }
3564
3565        let clock = get_atomic_clock_realtime();
3566        let mut emitter = ExecutionEventEmitter::new(
3567            clock,
3568            TraderId::from("TRADER-001"),
3569            AccountId::from("OKX-001"),
3570            AccountType::Margin,
3571            None,
3572        );
3573        let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel();
3574        emitter.set_sender(sender);
3575        let instruments = AtomicMap::new();
3576        let mut fee_cache = FeeCache::new();
3577        let mut filled_qty_cache = FilledQtyCache::new();
3578        let mut order_state_cache = AHashMap::new();
3579
3580        dispatch_ws_message(
3581            OKXWsMessage::SendFailed {
3582                request_id: "req-batch-send-failure".to_string(),
3583                client_order_ids: client_order_ids.to_vec(),
3584                op: Some(OKXWsOperation::BatchOrders),
3585                error,
3586            },
3587            &emitter,
3588            &state,
3589            AccountId::from("OKX-001"),
3590            AccountType::Margin,
3591            &instruments,
3592            &mut fee_cache,
3593            &mut filled_qty_cache,
3594            &mut order_state_cache,
3595            clock,
3596        );
3597
3598        for client_order_id in client_order_ids {
3599            assert_eq!(
3600                state.pending_orders.contains_key(client_order_id.as_str()),
3601                expected_pending,
3602                "pending state mismatch for {client_order_id}"
3603            );
3604        }
3605    }
3606}