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