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::Twap
840            | OKXAlgoOrderType::Chase
841            | OKXAlgoOrderType::Other
842    ) || msg.state == OKXAlgoOrderStatus::Unknown
843    {
844        return ExecutionUpdateRoute::Suppressed;
845    }
846
847    let direct_client_order_id = parse_client_order_id(&msg.algo_cl_ord_id)
848        .or_else(|| parse_client_order_id(&msg.cl_ord_id));
849
850    if let Some(client_order_id) = direct_client_order_id {
851        if state.contains_terminal(&client_order_id) {
852            return ExecutionUpdateRoute::Suppressed;
853        }
854
855        if let Some(context) = state
856            .order_contexts
857            .get(&client_order_id)
858            .map(|entry| *entry)
859        {
860            return ExecutionUpdateRoute::Tracked(client_order_id, context);
861        }
862
863        if state.order_identities.contains_key(&client_order_id) {
864            return ExecutionUpdateRoute::Suppressed;
865        }
866    }
867
868    let parent_venue_order_id = VenueOrderId::new(msg.algo_id.as_str());
869    let client_order_id = {
870        let bindings = state.lifecycle_bindings.lock();
871        bindings.client_order_id(&parent_venue_order_id)
872    };
873
874    let Some(client_order_id) = client_order_id else {
875        return ExecutionUpdateRoute::External;
876    };
877
878    if state.contains_terminal(&client_order_id) {
879        return ExecutionUpdateRoute::Suppressed;
880    }
881
882    state
883        .order_contexts
884        .get(&client_order_id)
885        .map_or(ExecutionUpdateRoute::Suppressed, |entry| {
886            ExecutionUpdateRoute::Tracked(client_order_id, *entry)
887        })
888}
889
890fn dispatch_algo_order_message(
891    msg: &OKXAlgoOrderMsg,
892    emitter: &ExecutionEventEmitter,
893    state: &WsDispatchState,
894    account_id: AccountId,
895    instruments: &AHashMap<Ustr, InstrumentAny>,
896    ts_init: UnixNanos,
897) {
898    let route = route_algo_order_message(msg, state);
899
900    match route {
901        ExecutionUpdateRoute::External => {
902            match parse_algo_order_msg(msg, account_id, instruments, ts_init) {
903                Ok(Some(report)) => dispatch_execution_reports(vec![report], emitter, state),
904                Ok(None) => {}
905                Err(e) => log::error!("Failed to parse external algo order message: {e}"),
906            }
907        }
908        ExecutionUpdateRoute::Suppressed => {
909            log::debug!(
910                "Suppressing algo order update: algo_id={} state={:?}",
911                msg.algo_id,
912                msg.state,
913            );
914        }
915        ExecutionUpdateRoute::Tracked(client_order_id, context) => {
916            let Some(instrument) = instruments.get(&msg.inst_id) else {
917                log::warn!(
918                    "No instrument for {}, skipping algo order message",
919                    msg.inst_id
920                );
921                return;
922            };
923            dispatch_tracked_algo_order_message(
924                msg,
925                client_order_id,
926                context,
927                instrument,
928                emitter,
929                state,
930                account_id,
931                ts_init,
932            );
933        }
934    }
935}
936
937#[expect(
938    clippy::too_many_arguments,
939    reason = "tracked routing requires the resolved context, venue state, and event timestamps"
940)]
941fn dispatch_tracked_algo_order_message(
942    msg: &OKXAlgoOrderMsg,
943    client_order_id: ClientOrderId,
944    mut context: OrderContext,
945    instrument: &InstrumentAny,
946    emitter: &ExecutionEventEmitter,
947    state: &WsDispatchState,
948    account_id: AccountId,
949    ts_init: UnixNanos,
950) {
951    let parent_venue_order_id = VenueOrderId::new(msg.algo_id.as_str());
952    let ts_event = parse_millisecond_timestamp(msg.u_time);
953    let mut bindings = state.lifecycle_bindings.lock();
954    let mut is_terminal = false;
955    let mut binding = bindings
956        .venue_by_client
957        .get(&client_order_id)
958        .copied()
959        .unwrap_or(OrderVenueBinding {
960            parent: parent_venue_order_id,
961            child: None,
962        });
963
964    if binding.parent != parent_venue_order_id {
965        log::error!(
966            "Suppressing conflicting algo parent binding for {client_order_id}: expected={} received={parent_venue_order_id}",
967            binding.parent,
968        );
969        return;
970    }
971
972    if binding.child.is_none() {
973        context = refresh_algo_order_context(msg, context, instrument, account_id, ts_init);
974        state.track_order_context(context);
975    }
976
977    bindings
978        .client_by_parent
979        .insert(parent_venue_order_id, client_order_id);
980
981    match msg.state {
982        OKXAlgoOrderStatus::Live | OKXAlgoOrderStatus::Pause => {
983            if binding.child.is_some() {
984                log::debug!(
985                    "Suppressing stale algo parent acceptance for {client_order_id}: algo_id={}",
986                    msg.algo_id,
987                );
988            } else {
989                ensure_accepted_emitted(
990                    client_order_id,
991                    account_id,
992                    parent_venue_order_id,
993                    &context.identity,
994                    emitter,
995                    state,
996                    ts_event,
997                    ts_init,
998                );
999            }
1000        }
1001        OKXAlgoOrderStatus::Effective
1002        | OKXAlgoOrderStatus::OrderPlaced
1003        | OKXAlgoOrderStatus::PartiallyEffective
1004        | OKXAlgoOrderStatus::Filled
1005        | OKXAlgoOrderStatus::PartiallyFailed => {
1006            let child_venue_order_id = algo_child_venue_order_id(msg);
1007            if let Some(child_venue_order_id) = child_venue_order_id {
1008                bind_algo_child_and_emit_transition(
1009                    client_order_id,
1010                    parent_venue_order_id,
1011                    child_venue_order_id,
1012                    &mut binding,
1013                    context,
1014                    account_id,
1015                    ts_event,
1016                    ts_init,
1017                    emitter,
1018                    state,
1019                );
1020            } else {
1021                ensure_accepted_emitted(
1022                    client_order_id,
1023                    account_id,
1024                    parent_venue_order_id,
1025                    &context.identity,
1026                    emitter,
1027                    state,
1028                    ts_event,
1029                    ts_init,
1030                );
1031            }
1032
1033            if matches!(
1034                msg.state,
1035                OKXAlgoOrderStatus::Filled | OKXAlgoOrderStatus::PartiallyFailed
1036            ) {
1037                log::debug!(
1038                    "Deferring tracked algo {:?} update for {client_order_id} to regular child execution updates",
1039                    msg.state,
1040                );
1041            }
1042        }
1043        OKXAlgoOrderStatus::Canceled => {
1044            if binding.child.is_some() {
1045                log::debug!(
1046                    "Suppressing stale canceled algo parent for triggered order {client_order_id}"
1047                );
1048            } else {
1049                ensure_accepted_emitted(
1050                    client_order_id,
1051                    account_id,
1052                    parent_venue_order_id,
1053                    &context.identity,
1054                    emitter,
1055                    state,
1056                    ts_event,
1057                    ts_init,
1058                );
1059                let canceled = OrderCanceled::new(
1060                    emitter.trader_id(),
1061                    context.identity.strategy_id,
1062                    context.identity.instrument_id,
1063                    client_order_id,
1064                    UUID4::new(),
1065                    ts_event,
1066                    ts_init,
1067                    false,
1068                    Some(parent_venue_order_id),
1069                    Some(account_id),
1070                );
1071                state.insert_terminal(client_order_id);
1072                state.remove_accepted(&client_order_id);
1073                state.remove_order_tracking(client_order_id);
1074                is_terminal = true;
1075                emitter.send_order_event(OrderEventAny::Canceled(canceled));
1076            }
1077        }
1078        OKXAlgoOrderStatus::OrderFailed => {
1079            if binding.child.is_some() {
1080                log::debug!(
1081                    "Suppressing stale failed algo parent for triggered order {client_order_id}"
1082                );
1083            } else {
1084                let reason = if msg.fail_code.is_empty() {
1085                    "OKX algo order failed"
1086                } else {
1087                    msg.fail_code.as_str()
1088                };
1089                let rejected = OrderRejected::new(
1090                    emitter.trader_id(),
1091                    context.identity.strategy_id,
1092                    context.identity.instrument_id,
1093                    client_order_id,
1094                    account_id,
1095                    Ustr::from(reason),
1096                    UUID4::new(),
1097                    ts_event,
1098                    ts_init,
1099                    false,
1100                    false,
1101                );
1102                state.insert_terminal(client_order_id);
1103                state.remove_accepted(&client_order_id);
1104                state.remove_order_tracking(client_order_id);
1105                is_terminal = true;
1106                emitter.send_order_event(OrderEventAny::Rejected(rejected));
1107            }
1108        }
1109        OKXAlgoOrderStatus::Unknown => {}
1110    }
1111
1112    if is_terminal {
1113        bindings.finish(client_order_id, binding);
1114    } else {
1115        bindings.venue_by_client.insert(client_order_id, binding);
1116    }
1117
1118    drop(bindings);
1119    state.linked_child_notify.notify_one();
1120}
1121
1122fn algo_child_venue_order_id(msg: &OKXAlgoOrderMsg) -> Option<VenueOrderId> {
1123    if !msg.ord_id.is_empty() {
1124        return Some(VenueOrderId::new(msg.ord_id.as_str()));
1125    }
1126
1127    match msg.ord_id_list.as_slice() {
1128        [order_id] if !order_id.is_empty() => Some(VenueOrderId::new(order_id.as_str())),
1129        [] => None,
1130        order_ids => {
1131            log::warn!(
1132                "Cannot bind algo order {} to {} triggered child IDs",
1133                msg.algo_id,
1134                order_ids.len(),
1135            );
1136            None
1137        }
1138    }
1139}
1140
1141fn refresh_algo_order_context(
1142    msg: &OKXAlgoOrderMsg,
1143    mut context: OrderContext,
1144    instrument: &InstrumentAny,
1145    account_id: AccountId,
1146    ts_init: UnixNanos,
1147) -> OrderContext {
1148    match parse_algo_order_status_report(msg, instrument, account_id, ts_init) {
1149        Ok(report) => {
1150            if !msg.sz.is_empty() {
1151                context.quantity = report.quantity;
1152            }
1153
1154            context.price = report.price;
1155            context.trigger_price = report.trigger_price;
1156            context.trigger_type = report.trigger_type;
1157        }
1158        Err(e) => {
1159            log::error!(
1160                "Failed to refresh tracked algo order context for {}: {e}",
1161                context.identity.client_order_id,
1162            );
1163            return context;
1164        }
1165    }
1166
1167    if !msg.actual_sz.is_empty() && msg.actual_sz != "0" {
1168        match parse_quantity(msg.actual_sz.as_str(), instrument.size_precision()) {
1169            Ok(quantity) => context.quantity = quantity,
1170            Err(e) => log::error!(
1171                "Failed to refresh tracked algo actual quantity for {}: {e}",
1172                context.identity.client_order_id,
1173            ),
1174        }
1175    }
1176
1177    context
1178}
1179
1180fn refresh_regular_child_context(
1181    msg: &OKXOrderMsg,
1182    mut context: OrderContext,
1183    instrument: &InstrumentAny,
1184) -> OrderContext {
1185    match parse_quantity(&msg.sz, instrument.size_precision()) {
1186        Ok(quantity) => context.quantity = quantity,
1187        Err(e) => log::error!(
1188            "Failed to refresh tracked child quantity for {}: {e}",
1189            context.identity.client_order_id,
1190        ),
1191    }
1192
1193    context.price = if is_market_price(&msg.px) {
1194        None
1195    } else {
1196        match parse_price(&msg.px, instrument.price_precision()) {
1197            Ok(price) => Some(price),
1198            Err(e) => {
1199                log::error!(
1200                    "Failed to refresh tracked child price for {}: {e}",
1201                    context.identity.client_order_id,
1202                );
1203                context.price
1204            }
1205        }
1206    };
1207    context
1208}
1209
1210#[expect(clippy::too_many_arguments)]
1211fn bind_algo_child_and_emit_transition(
1212    client_order_id: ClientOrderId,
1213    parent_venue_order_id: VenueOrderId,
1214    child_venue_order_id: VenueOrderId,
1215    binding: &mut OrderVenueBinding,
1216    context: OrderContext,
1217    account_id: AccountId,
1218    ts_event: UnixNanos,
1219    ts_init: UnixNanos,
1220    emitter: &ExecutionEventEmitter,
1221    state: &WsDispatchState,
1222) {
1223    if let Some(bound_child) = binding.child {
1224        if bound_child != child_venue_order_id {
1225            log::error!(
1226                "Suppressing conflicting algo child binding for {client_order_id}: expected={bound_child} received={child_venue_order_id}"
1227            );
1228        }
1229        return;
1230    }
1231
1232    binding.child = Some(child_venue_order_id);
1233    state.track_order_context(context);
1234    ensure_accepted_emitted(
1235        client_order_id,
1236        account_id,
1237        parent_venue_order_id,
1238        &context.identity,
1239        emitter,
1240        state,
1241        ts_event,
1242        ts_init,
1243    );
1244
1245    if state.accepted_venue_order_id(&client_order_id) != Some(child_venue_order_id) {
1246        state.insert_accepted(client_order_id, child_venue_order_id);
1247        emit_child_update(
1248            client_order_id,
1249            child_venue_order_id,
1250            context,
1251            account_id,
1252            ts_event,
1253            ts_init,
1254            emitter,
1255        );
1256    }
1257
1258    if !state.contains_triggered(&client_order_id) {
1259        state.insert_triggered(client_order_id);
1260
1261        if TRIGGERABLE_ORDER_TYPES.contains(&context.identity.order_type) {
1262            let triggered = OrderTriggered::new(
1263                emitter.trader_id(),
1264                context.identity.strategy_id,
1265                context.identity.instrument_id,
1266                client_order_id,
1267                UUID4::new(),
1268                ts_event,
1269                ts_init,
1270                false,
1271                Some(child_venue_order_id),
1272                Some(account_id),
1273            );
1274            emitter.send_order_event(OrderEventAny::Triggered(triggered));
1275        }
1276    }
1277}
1278
1279#[expect(clippy::too_many_arguments)]
1280fn refresh_bound_child_and_emit_update(
1281    client_order_id: ClientOrderId,
1282    child_venue_order_id: VenueOrderId,
1283    context: OrderContext,
1284    account_id: AccountId,
1285    ts_event: UnixNanos,
1286    ts_init: UnixNanos,
1287    emitter: &ExecutionEventEmitter,
1288    state: &WsDispatchState,
1289    emit_update: bool,
1290) {
1291    let terms_changed = state
1292        .order_contexts
1293        .get(&client_order_id)
1294        .is_some_and(|previous| {
1295            previous.quantity != context.quantity
1296                || previous.price != context.price
1297                || previous.trigger_price != context.trigger_price
1298        });
1299    state.track_order_context(context);
1300
1301    if !terms_changed || !emit_update {
1302        return;
1303    }
1304
1305    state.insert_accepted(client_order_id, child_venue_order_id);
1306    emit_child_update(
1307        client_order_id,
1308        child_venue_order_id,
1309        context,
1310        account_id,
1311        ts_event,
1312        ts_init,
1313        emitter,
1314    );
1315}
1316
1317fn emit_child_update(
1318    client_order_id: ClientOrderId,
1319    child_venue_order_id: VenueOrderId,
1320    context: OrderContext,
1321    account_id: AccountId,
1322    ts_event: UnixNanos,
1323    ts_init: UnixNanos,
1324    emitter: &ExecutionEventEmitter,
1325) {
1326    let updated = OrderUpdated::new(
1327        emitter.trader_id(),
1328        context.identity.strategy_id,
1329        context.identity.instrument_id,
1330        client_order_id,
1331        context.quantity,
1332        UUID4::new(),
1333        ts_event,
1334        ts_init,
1335        false,
1336        Some(child_venue_order_id),
1337        Some(account_id),
1338        context.price,
1339        context.trigger_price,
1340        None,
1341        false,
1342    );
1343    emitter.send_order_event(OrderEventAny::Updated(updated));
1344}
1345
1346/// Dispatches order messages, producing proper order events for tracked orders
1347/// and falling back to execution reports for untracked/external orders.
1348#[expect(clippy::too_many_arguments)]
1349fn dispatch_order_messages(
1350    order_msgs: &[OKXOrderMsg],
1351    emitter: &ExecutionEventEmitter,
1352    state: &WsDispatchState,
1353    account_id: AccountId,
1354    instruments: &AHashMap<Ustr, InstrumentAny>,
1355    fee_cache: &mut AHashMap<Ustr, Money>,
1356    filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
1357    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
1358    ts_init: UnixNanos,
1359) {
1360    for msg in order_msgs {
1361        let Some(instrument) = instruments.get(&msg.inst_id) else {
1362            log::warn!("No instrument for {}, skipping order message", msg.inst_id);
1363            continue;
1364        };
1365
1366        let direct_client_order_id = parse_client_order_id(&msg.cl_ord_id);
1367        let parent_client_order_id = msg
1368            .algo_cl_ord_id
1369            .as_deref()
1370            .and_then(parse_client_order_id);
1371        let linked_parent_venue_order_id = msg
1372            .algo_id
1373            .as_deref()
1374            .filter(|value| !value.is_empty())
1375            .or_else(|| {
1376                msg.linked_algo_ord
1377                    .as_ref()
1378                    .map(|linked| linked.algo_id.as_str())
1379                    .filter(|value| !value.is_empty())
1380            })
1381            .map(VenueOrderId::new);
1382
1383        // Triggered child orders may have a generated or empty cl_ord_id.
1384        // Resolve the tracked parent before falling back to a report.
1385        let direct_resolution = [direct_client_order_id, parent_client_order_id]
1386            .into_iter()
1387            .flatten()
1388            .find_map(|client_order_id| {
1389                state
1390                    .order_identity(client_order_id)
1391                    .map(|identity| (client_order_id, Some(identity)))
1392            });
1393
1394        let bound_client_order_id = if direct_resolution.is_none()
1395            && let Some(parent_venue_order_id) = linked_parent_venue_order_id
1396        {
1397            match state.resolve_or_hold_linked_child(parent_venue_order_id, instrument.id(), msg) {
1398                LinkedChildResolution::Bound(client_order_id) => Some(client_order_id),
1399                LinkedChildResolution::Held => {
1400                    log::debug!(
1401                        "Holding linked child update during algo parent binding: ord_id={} parent_id={parent_venue_order_id}",
1402                        msg.ord_id,
1403                    );
1404                    continue;
1405                }
1406                LinkedChildResolution::External => None,
1407            }
1408        } else {
1409            None
1410        };
1411
1412        let resolved = direct_resolution
1413            .or_else(|| {
1414                bound_client_order_id.and_then(|client_order_id| {
1415                    state
1416                        .order_identity(client_order_id)
1417                        .map(|identity| (client_order_id, Some(identity)))
1418                })
1419            })
1420            .or_else(|| {
1421                bound_client_order_id
1422                    .or(parent_client_order_id)
1423                    .or(direct_client_order_id)
1424                    .map(|client_order_id| (client_order_id, None))
1425            });
1426
1427        let Some((client_order_id, identity)) = resolved else {
1428            log::debug!(
1429                "Order without client or algo client order ID (ord_id={}), sending as report",
1430                msg.ord_id
1431            );
1432            dispatch_order_msg_as_report(
1433                msg,
1434                account_id,
1435                instruments,
1436                fee_cache,
1437                filled_qty_cache,
1438                emitter,
1439                state,
1440                ts_init,
1441            );
1442            continue;
1443        };
1444
1445        if let Some(ident) = identity {
1446            let context = state
1447                .order_contexts
1448                .get(&client_order_id)
1449                .map(|entry| refresh_regular_child_context(msg, *entry, instrument));
1450            let mut lifecycle_bindings = context.map(|_| state.lifecycle_bindings.lock());
1451
1452            if let (Some(context), Some(bindings)) = (context, lifecycle_bindings.as_mut()) {
1453                let parent_venue_order_id = linked_parent_venue_order_id.or_else(|| {
1454                    bindings
1455                        .venue_by_client
1456                        .get(&client_order_id)
1457                        .map(|binding| binding.parent)
1458                });
1459
1460                let child_venue_order_id = VenueOrderId::new(msg.ord_id);
1461                let parent_venue_order_id = parent_venue_order_id.unwrap_or(child_venue_order_id);
1462                let mut binding = bindings
1463                    .venue_by_client
1464                    .get(&client_order_id)
1465                    .copied()
1466                    .unwrap_or(OrderVenueBinding {
1467                        parent: parent_venue_order_id,
1468                        child: None,
1469                    });
1470
1471                if binding.parent != parent_venue_order_id {
1472                    log::error!(
1473                        "Suppressing conflicting child parent binding for {client_order_id}: expected={} received={parent_venue_order_id}",
1474                        binding.parent,
1475                    );
1476                    continue;
1477                }
1478
1479                bindings
1480                    .client_by_parent
1481                    .insert(parent_venue_order_id, client_order_id);
1482                let ts_event = parse_millisecond_timestamp(msg.u_time);
1483
1484                if binding.child == Some(child_venue_order_id) {
1485                    refresh_bound_child_and_emit_update(
1486                        client_order_id,
1487                        child_venue_order_id,
1488                        context,
1489                        account_id,
1490                        ts_event,
1491                        ts_init,
1492                        emitter,
1493                        state,
1494                        !order_state_cache.contains_key(&client_order_id),
1495                    );
1496                } else {
1497                    bind_algo_child_and_emit_transition(
1498                        client_order_id,
1499                        parent_venue_order_id,
1500                        child_venue_order_id,
1501                        &mut binding,
1502                        context,
1503                        account_id,
1504                        ts_event,
1505                        ts_init,
1506                        emitter,
1507                        state,
1508                    );
1509                }
1510
1511                bindings.venue_by_client.insert(client_order_id, binding);
1512            }
1513
1514            let is_post_only_cancel = is_post_only_auto_cancel(msg);
1515
1516            if is_post_only_cancel
1517                || (!state.contains_accepted(&client_order_id) && is_unfilled_rpi_cancel(msg))
1518            {
1519                if is_post_only_cancel {
1520                    state.insert_post_only_rejection(msg.ord_id);
1521                }
1522
1523                let ts_event = parse_millisecond_timestamp(msg.u_time);
1524                let reason = if msg.ord_type == OKXOrderType::Rpi {
1525                    msg.cancel_source_reason
1526                        .as_deref()
1527                        .filter(|reason| !reason.is_empty())
1528                        .unwrap_or("RPI order canceled before acceptance")
1529                } else {
1530                    "Post-only order would have taken liquidity"
1531                };
1532                let rejected = OrderRejected::new(
1533                    emitter.trader_id(),
1534                    ident.strategy_id,
1535                    instrument.id(),
1536                    client_order_id,
1537                    account_id,
1538                    Ustr::from(reason),
1539                    UUID4::new(),
1540                    ts_event,
1541                    ts_init,
1542                    false,
1543                    true, // due_post_only
1544                );
1545                state.remove_order_tracking(client_order_id);
1546                if let Some(bindings) = lifecycle_bindings.as_mut() {
1547                    state.insert_terminal(client_order_id);
1548                    if let Some(binding) = bindings.venue_by_client.get(&client_order_id).copied() {
1549                        bindings.finish(client_order_id, binding);
1550                    }
1551                }
1552
1553                order_state_cache.remove(&client_order_id);
1554                fee_cache.remove(&msg.ord_id);
1555                filled_qty_cache.remove(&msg.ord_id);
1556                emitter.send_order_event(OrderEventAny::Rejected(rejected));
1557                continue;
1558            }
1559
1560            let previous_fee = fee_cache.get(&msg.ord_id).copied();
1561            let previous_filled_qty = filled_qty_cache.get(&msg.ord_id).copied();
1562            let previous_state = order_state_cache.get(&client_order_id);
1563
1564            match parse_order_event(
1565                msg,
1566                client_order_id,
1567                account_id,
1568                emitter.trader_id(),
1569                ident.strategy_id,
1570                instrument,
1571                previous_fee,
1572                previous_filled_qty,
1573                previous_state,
1574                ts_init,
1575            ) {
1576                Ok(event) => {
1577                    update_order_state_cache(msg, instrument, client_order_id, order_state_cache);
1578                    dispatch_parsed_order_event(
1579                        event,
1580                        client_order_id,
1581                        account_id,
1582                        VenueOrderId::new(msg.ord_id),
1583                        &ident,
1584                        instrument,
1585                        msg.state,
1586                        emitter,
1587                        state,
1588                        order_state_cache,
1589                        ts_init,
1590                    );
1591
1592                    if state.contains_terminal(&client_order_id)
1593                        && let Some(bindings) = lifecycle_bindings.as_mut()
1594                        && let Some(binding) =
1595                            bindings.venue_by_client.get(&client_order_id).copied()
1596                    {
1597                        bindings.finish(client_order_id, binding);
1598                    }
1599
1600                    update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
1601                }
1602                Err(e) => log::error!("Failed to parse order event for {client_order_id}: {e}"),
1603            }
1604        } else if state.contains_terminal(&client_order_id) {
1605            dispatch_terminal_order_fill_as_report(
1606                msg,
1607                client_order_id,
1608                account_id,
1609                instruments,
1610                fee_cache,
1611                filled_qty_cache,
1612                emitter,
1613                state,
1614                ts_init,
1615            );
1616        } else if is_post_only_auto_cancel(msg) && state.contains_post_only_rejection(&msg.ord_id) {
1617            log::debug!(
1618                "Skipping replayed post-only rejection for {client_order_id}: ord_id={}",
1619                msg.ord_id
1620            );
1621        } else {
1622            log::debug!(
1623                "Untracked order {client_order_id} (ord_id={}), sending as report for reconciliation",
1624                msg.ord_id
1625            );
1626            dispatch_order_msg_as_report(
1627                msg,
1628                account_id,
1629                instruments,
1630                fee_cache,
1631                filled_qty_cache,
1632                emitter,
1633                state,
1634                ts_init,
1635            );
1636        }
1637    }
1638}
1639
1640#[expect(clippy::too_many_arguments)]
1641fn dispatch_spread_order_messages(
1642    order_msgs: &[OKXSpreadOrder],
1643    emitter: &ExecutionEventEmitter,
1644    state: &WsDispatchState,
1645    account_id: AccountId,
1646    instruments: &AHashMap<Ustr, InstrumentAny>,
1647    filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
1648    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
1649    ts_init: UnixNanos,
1650) {
1651    for msg in order_msgs {
1652        let Some(instrument) = instruments.get(&msg.sprd_id) else {
1653            log::warn!(
1654                "No instrument for {}, skipping spread order message",
1655                msg.sprd_id
1656            );
1657            continue;
1658        };
1659
1660        let Some(client_order_id) = parse_client_order_id(msg.cl_ord_id.as_str()) else {
1661            log::debug!(
1662                "Spread order without client_order_id (ord_id={}), sending as report",
1663                msg.ord_id
1664            );
1665            dispatch_spread_order_msg_as_report(
1666                msg,
1667                account_id,
1668                instruments,
1669                filled_qty_cache,
1670                emitter,
1671                state,
1672                ts_init,
1673            );
1674            continue;
1675        };
1676
1677        let identity = state.order_identity(client_order_id);
1678
1679        if let Some(ident) = identity {
1680            if is_spread_post_only_auto_cancel(msg) {
1681                let ts_event = msg
1682                    .u_time
1683                    .or(msg.c_time)
1684                    .map_or(ts_init, parse_millisecond_timestamp);
1685                let rejected = OrderRejected::new(
1686                    emitter.trader_id(),
1687                    ident.strategy_id,
1688                    instrument.id(),
1689                    client_order_id,
1690                    account_id,
1691                    Ustr::from(OKX_POST_ONLY_CANCEL_REASON),
1692                    UUID4::new(),
1693                    ts_event,
1694                    ts_init,
1695                    false,
1696                    true,
1697                );
1698                state.remove_order_tracking(client_order_id);
1699                order_state_cache.remove(&client_order_id);
1700                filled_qty_cache.remove(&msg.ord_id);
1701                emitter.send_order_event(OrderEventAny::Rejected(rejected));
1702                continue;
1703            }
1704
1705            let previous_filled_qty = filled_qty_cache.get(&msg.ord_id).copied();
1706            let previous_state = order_state_cache.get(&client_order_id);
1707
1708            match parse_spread_order_event(
1709                msg,
1710                client_order_id,
1711                account_id,
1712                emitter.trader_id(),
1713                ident.strategy_id,
1714                instrument,
1715                previous_filled_qty,
1716                previous_state,
1717                ts_init,
1718            ) {
1719                Ok(event) => {
1720                    update_spread_order_state_cache(
1721                        msg,
1722                        instrument,
1723                        client_order_id,
1724                        order_state_cache,
1725                    );
1726                    dispatch_parsed_order_event(
1727                        event,
1728                        client_order_id,
1729                        account_id,
1730                        VenueOrderId::new(msg.ord_id.as_str()),
1731                        &ident,
1732                        instrument,
1733                        msg.state,
1734                        emitter,
1735                        state,
1736                        order_state_cache,
1737                        ts_init,
1738                    );
1739                    update_spread_fill_cache(msg, instrument, filled_qty_cache);
1740                }
1741                Err(e) => {
1742                    log::error!("Failed to parse spread order event for {client_order_id}: {e}");
1743                }
1744            }
1745        } else {
1746            log::debug!(
1747                "Untracked spread order {client_order_id} (ord_id={}), sending as report for reconciliation",
1748                msg.ord_id
1749            );
1750            dispatch_spread_order_msg_as_report(
1751                msg,
1752                account_id,
1753                instruments,
1754                filled_qty_cache,
1755                emitter,
1756                state,
1757                ts_init,
1758            );
1759        }
1760    }
1761}
1762
1763/// Dispatches a parsed order event as a proper `OrderEventAny`.
1764///
1765/// Guarantees the `Submitted -> Accepted -> ...` lifecycle by synthesizing
1766/// `OrderAccepted` before any other event when one has not yet been emitted.
1767/// Duplicate `Accepted` events (e.g. from reconnect replays) are suppressed.
1768#[expect(clippy::too_many_arguments)]
1769fn dispatch_parsed_order_event(
1770    event: ParsedOrderEvent,
1771    client_order_id: ClientOrderId,
1772    account_id: AccountId,
1773    venue_order_id: VenueOrderId,
1774    identity: &OrderIdentity,
1775    instrument: &InstrumentAny,
1776    venue_status: OKXOrderStatus,
1777    emitter: &ExecutionEventEmitter,
1778    state: &WsDispatchState,
1779    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
1780    ts_init: UnixNanos,
1781) {
1782    let is_terminal;
1783
1784    match event {
1785        ParsedOrderEvent::Accepted(e) => {
1786            if state.contains_filled(&client_order_id) || state.contains_terminal(&client_order_id)
1787            {
1788                log::debug!("Skipping duplicate Accepted for {client_order_id}");
1789                return;
1790            }
1791
1792            if state.contains_accepted(&client_order_id) {
1793                emit_venue_order_id_update_if_changed(
1794                    client_order_id,
1795                    account_id,
1796                    venue_order_id,
1797                    identity,
1798                    e.ts_event,
1799                    emitter,
1800                    state,
1801                    order_state_cache,
1802                    ts_init,
1803                );
1804                return;
1805            }
1806
1807            if state.contains_triggered(&client_order_id) {
1808                log::debug!("Skipping duplicate Accepted for {client_order_id}");
1809                return;
1810            }
1811
1812            state.insert_accepted(client_order_id, venue_order_id);
1813            is_terminal = false;
1814            emitter.send_order_event(OrderEventAny::Accepted(e));
1815        }
1816        ParsedOrderEvent::Triggered(e) => {
1817            if state.contains_filled(&client_order_id) {
1818                log::debug!("Skipping stale Triggered for {client_order_id} (already filled)");
1819                return;
1820            }
1821
1822            if !TRIGGERABLE_ORDER_TYPES.contains(&identity.order_type) {
1823                log::debug!(
1824                    "Skipping OrderTriggered for {} order {client_order_id}: market-style stops have no TRIGGERED state",
1825                    identity.order_type,
1826                );
1827                state.insert_triggered(client_order_id);
1828                return;
1829            }
1830
1831            ensure_accepted_emitted(
1832                client_order_id,
1833                account_id,
1834                venue_order_id,
1835                identity,
1836                emitter,
1837                state,
1838                ts_init,
1839                ts_init,
1840            );
1841            state.insert_triggered(client_order_id);
1842            is_terminal = false;
1843            emitter.send_order_event(OrderEventAny::Triggered(e));
1844        }
1845        ParsedOrderEvent::Canceled(e) => {
1846            ensure_accepted_emitted(
1847                client_order_id,
1848                account_id,
1849                venue_order_id,
1850                identity,
1851                emitter,
1852                state,
1853                ts_init,
1854                ts_init,
1855            );
1856            state.remove_triggered(&client_order_id);
1857            state.remove_filled(&client_order_id);
1858            is_terminal = true;
1859            emitter.send_order_event(OrderEventAny::Canceled(e));
1860        }
1861        ParsedOrderEvent::Expired(e) => {
1862            ensure_accepted_emitted(
1863                client_order_id,
1864                account_id,
1865                venue_order_id,
1866                identity,
1867                emitter,
1868                state,
1869                ts_init,
1870                ts_init,
1871            );
1872            state.remove_triggered(&client_order_id);
1873            state.remove_filled(&client_order_id);
1874            is_terminal = true;
1875            emitter.send_order_event(OrderEventAny::Expired(e));
1876        }
1877        ParsedOrderEvent::Updated(e) => {
1878            ensure_accepted_emitted(
1879                client_order_id,
1880                account_id,
1881                venue_order_id,
1882                identity,
1883                emitter,
1884                state,
1885                ts_init,
1886                ts_init,
1887            );
1888            is_terminal = false;
1889            emitter.send_order_event(OrderEventAny::Updated(e));
1890        }
1891        ParsedOrderEvent::Fill(fill_report) => {
1892            is_terminal = venue_status == OKXOrderStatus::Filled;
1893
1894            if state.check_and_insert_trade(fill_report.trade_id) {
1895                log::debug!(
1896                    "Skipping duplicate fill for {client_order_id}: trade_id={}",
1897                    fill_report.trade_id
1898                );
1899            } else {
1900                emit_venue_order_id_update_if_changed(
1901                    client_order_id,
1902                    account_id,
1903                    venue_order_id,
1904                    identity,
1905                    fill_report.ts_event,
1906                    emitter,
1907                    state,
1908                    order_state_cache,
1909                    ts_init,
1910                );
1911                ensure_accepted_emitted(
1912                    client_order_id,
1913                    account_id,
1914                    venue_order_id,
1915                    identity,
1916                    emitter,
1917                    state,
1918                    ts_init,
1919                    ts_init,
1920                );
1921                state.insert_filled(client_order_id);
1922                state.remove_triggered(&client_order_id);
1923                let filled = fill_report_to_order_filled(
1924                    &fill_report,
1925                    emitter.trader_id(),
1926                    identity,
1927                    instrument.quote_currency(),
1928                );
1929                emitter.send_order_event(OrderEventAny::Filled(filled));
1930            }
1931        }
1932        ParsedOrderEvent::StatusOnly(report) => {
1933            is_terminal = matches!(
1934                report.order_status,
1935                OrderStatus::Filled | OrderStatus::Canceled | OrderStatus::Expired
1936            );
1937            emitter.send_order_status_report(*report);
1938        }
1939        ParsedOrderEvent::Skipped => return,
1940    }
1941
1942    if is_terminal {
1943        state.insert_terminal(client_order_id);
1944        state.remove_order_tracking(client_order_id);
1945        state.remove_accepted(&client_order_id);
1946        order_state_cache.remove(&client_order_id);
1947        // Keep fee_cache and filled_qty_cache entries: replayed terminal
1948        // messages go through the untracked report path and need prior
1949        // cumulative state to avoid re-emitting the full fill quantity
1950    }
1951}
1952
1953/// Synthesizes and emits `OrderAccepted` if one has not yet been emitted for
1954/// this order. Handles fast-filling orders that skip the `Live` state on OKX.
1955#[expect(
1956    clippy::too_many_arguments,
1957    reason = "acceptance reconstruction requires identity, venue state, and event timestamps"
1958)]
1959fn ensure_accepted_emitted(
1960    client_order_id: ClientOrderId,
1961    account_id: AccountId,
1962    venue_order_id: VenueOrderId,
1963    identity: &OrderIdentity,
1964    emitter: &ExecutionEventEmitter,
1965    state: &WsDispatchState,
1966    ts_event: UnixNanos,
1967    ts_init: UnixNanos,
1968) {
1969    if state.contains_accepted(&client_order_id)
1970        || state.contains_terminal(&client_order_id)
1971        || state.contains_filled(&client_order_id)
1972    {
1973        return;
1974    }
1975
1976    state.insert_accepted(client_order_id, venue_order_id);
1977    let accepted = OrderAccepted::new(
1978        emitter.trader_id(),
1979        identity.strategy_id,
1980        identity.instrument_id,
1981        client_order_id,
1982        venue_order_id,
1983        account_id,
1984        UUID4::new(),
1985        ts_event,
1986        ts_init,
1987        false,
1988    );
1989    emitter.send_order_event(OrderEventAny::Accepted(accepted));
1990}
1991
1992#[expect(clippy::too_many_arguments)]
1993fn emit_venue_order_id_update_if_changed(
1994    client_order_id: ClientOrderId,
1995    account_id: AccountId,
1996    venue_order_id: VenueOrderId,
1997    identity: &OrderIdentity,
1998    ts_event: UnixNanos,
1999    emitter: &ExecutionEventEmitter,
2000    state: &WsDispatchState,
2001    order_state_cache: &AHashMap<ClientOrderId, OrderStateSnapshot>,
2002    ts_init: UnixNanos,
2003) {
2004    let Some(accepted_venue_order_id) = state.accepted_venue_order_id(&client_order_id) else {
2005        return;
2006    };
2007
2008    if accepted_venue_order_id == venue_order_id {
2009        return;
2010    }
2011
2012    let Some(snapshot) = order_state_cache.get(&client_order_id) else {
2013        return;
2014    };
2015
2016    state.insert_accepted(client_order_id, venue_order_id);
2017    let updated = OrderUpdated::new(
2018        emitter.trader_id(),
2019        identity.strategy_id,
2020        identity.instrument_id,
2021        client_order_id,
2022        snapshot.quantity,
2023        UUID4::new(),
2024        ts_event,
2025        ts_init,
2026        false,
2027        Some(venue_order_id),
2028        Some(account_id),
2029        snapshot.price,
2030        None,
2031        None,
2032        false,
2033    );
2034    emitter.send_order_event(OrderEventAny::Updated(updated));
2035}
2036
2037/// Converts a [`FillReport`] into an [`OrderFilled`] event using tracked identity.
2038fn fill_report_to_order_filled(
2039    report: &FillReport,
2040    trader_id: TraderId,
2041    identity: &OrderIdentity,
2042    quote_currency: Currency,
2043) -> OrderFilled {
2044    OrderFilled::new(
2045        trader_id,
2046        identity.strategy_id,
2047        report.instrument_id,
2048        report
2049            .client_order_id
2050            .expect("tracked order has client_order_id"),
2051        report.venue_order_id,
2052        report.account_id,
2053        report.trade_id,
2054        identity.order_side,
2055        identity.order_type,
2056        report.last_qty,
2057        report.last_px,
2058        quote_currency,
2059        report.liquidity_side,
2060        UUID4::new(),
2061        report.ts_event,
2062        report.ts_init,
2063        false,
2064        report.venue_position_id,
2065        Some(report.commission),
2066        None,
2067    )
2068}
2069
2070/// Falls back to the report path for a single order message.
2071#[expect(clippy::too_many_arguments)]
2072fn dispatch_order_msg_as_report(
2073    msg: &OKXOrderMsg,
2074    account_id: AccountId,
2075    instruments: &AHashMap<Ustr, InstrumentAny>,
2076    fee_cache: &mut AHashMap<Ustr, Money>,
2077    filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
2078    emitter: &ExecutionEventEmitter,
2079    state: &WsDispatchState,
2080    ts_init: UnixNanos,
2081) {
2082    match parse_order_msg(
2083        msg,
2084        account_id,
2085        instruments,
2086        fee_cache,
2087        filled_qty_cache,
2088        ts_init,
2089    ) {
2090        Ok(report) => {
2091            dispatch_execution_reports(vec![report], emitter, state);
2092
2093            if let Some(instrument) = instruments.get(&msg.inst_id) {
2094                update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
2095            }
2096        }
2097        Err(e) => log::error!("Failed to parse order message as report: {e}"),
2098    }
2099}
2100
2101#[expect(clippy::too_many_arguments)]
2102fn dispatch_terminal_order_fill_as_report(
2103    msg: &OKXOrderMsg,
2104    client_order_id: ClientOrderId,
2105    account_id: AccountId,
2106    instruments: &AHashMap<Ustr, InstrumentAny>,
2107    fee_cache: &mut AHashMap<Ustr, Money>,
2108    filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
2109    emitter: &ExecutionEventEmitter,
2110    state: &WsDispatchState,
2111    ts_init: UnixNanos,
2112) {
2113    match parse_order_msg(
2114        msg,
2115        account_id,
2116        instruments,
2117        fee_cache,
2118        filled_qty_cache,
2119        ts_init,
2120    ) {
2121        Ok(ExecutionReport::Fill(mut report)) => {
2122            report.client_order_id = Some(client_order_id);
2123            dispatch_execution_reports(vec![ExecutionReport::Fill(report)], emitter, state);
2124
2125            if let Some(instrument) = instruments.get(&msg.inst_id) {
2126                update_fee_fill_caches(msg, instrument, fee_cache, filled_qty_cache);
2127            }
2128        }
2129        Ok(ExecutionReport::Order(_)) => {
2130            log::debug!(
2131                "Suppressing stale regular order status for terminal tracked order {client_order_id}: ord_id={}",
2132                msg.ord_id,
2133            );
2134        }
2135        Err(e) => log::error!("Failed to parse terminal order update: {e}"),
2136    }
2137}
2138
2139fn dispatch_spread_order_msg_as_report(
2140    msg: &OKXSpreadOrder,
2141    account_id: AccountId,
2142    instruments: &AHashMap<Ustr, InstrumentAny>,
2143    filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
2144    emitter: &ExecutionEventEmitter,
2145    state: &WsDispatchState,
2146    ts_init: UnixNanos,
2147) {
2148    match parse_spread_order_msg(msg, account_id, instruments, filled_qty_cache, ts_init) {
2149        Ok(report) => {
2150            dispatch_execution_reports(vec![report], emitter, state);
2151
2152            if let Some(instrument) = instruments.get(&msg.sprd_id) {
2153                update_spread_fill_cache(msg, instrument, filled_qty_cache);
2154            }
2155        }
2156        Err(e) => log::error!("Failed to parse spread order message as report: {e}"),
2157    }
2158}
2159
2160/// Updates fee, fill, and order state caches from a raw OKX order message.
2161fn update_order_state_cache(
2162    msg: &OKXOrderMsg,
2163    instrument: &InstrumentAny,
2164    client_order_id: ClientOrderId,
2165    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
2166) {
2167    let venue_order_id = VenueOrderId::new(msg.ord_id);
2168    let quantity = parse_quantity(&msg.sz, instrument.size_precision()).unwrap_or_default();
2169    let price = if is_market_price(&msg.px) {
2170        None
2171    } else {
2172        parse_price(&msg.px, instrument.price_precision()).ok()
2173    };
2174
2175    order_state_cache.insert(
2176        client_order_id,
2177        OrderStateSnapshot {
2178            venue_order_id,
2179            quantity,
2180            price,
2181        },
2182    );
2183}
2184
2185fn update_spread_order_state_cache(
2186    msg: &OKXSpreadOrder,
2187    instrument: &InstrumentAny,
2188    client_order_id: ClientOrderId,
2189    order_state_cache: &mut AHashMap<ClientOrderId, OrderStateSnapshot>,
2190) {
2191    let venue_order_id = VenueOrderId::new(msg.ord_id.as_str());
2192    let quantity = parse_quantity(&msg.sz, instrument.size_precision()).unwrap_or_default();
2193    let price = if is_market_price(&msg.px) {
2194        None
2195    } else {
2196        parse_price(&msg.px, instrument.price_precision()).ok()
2197    };
2198
2199    order_state_cache.insert(
2200        client_order_id,
2201        OrderStateSnapshot {
2202            venue_order_id,
2203            quantity,
2204            price,
2205        },
2206    );
2207}
2208
2209fn update_spread_fill_cache(
2210    msg: &OKXSpreadOrder,
2211    instrument: &InstrumentAny,
2212    filled_qty_cache: &mut AHashMap<Ustr, Quantity>,
2213) {
2214    if !msg.acc_fill_sz.is_empty()
2215        && msg.acc_fill_sz != "0"
2216        && let Ok(qty) = parse_quantity(&msg.acc_fill_sz, instrument.size_precision())
2217    {
2218        filled_qty_cache.insert(msg.ord_id, qty);
2219    }
2220}
2221
2222fn is_spread_post_only_auto_cancel(msg: &OKXSpreadOrder) -> bool {
2223    msg.state == OKXOrderStatus::Canceled && msg.cancel_source == OKX_POST_ONLY_CANCEL_SOURCE
2224}
2225
2226/// Dispatches execution reports with cross-stream deduplication.
2227pub fn dispatch_execution_reports(
2228    reports: Vec<ExecutionReport>,
2229    emitter: &ExecutionEventEmitter,
2230    state: &WsDispatchState,
2231) {
2232    log::debug!("Processing {} execution report(s)", reports.len());
2233
2234    for report in reports {
2235        match report {
2236            ExecutionReport::Order(order_report) => {
2237                if let Some(cid) = order_report.client_order_id {
2238                    match order_report.order_status {
2239                        // Guard form reformats awkwardly across multiple lines
2240                        #[allow(clippy::collapsible_match)]
2241                        OrderStatus::Accepted => {
2242                            if state.contains_terminal(&cid)
2243                                || state.contains_filled(&cid)
2244                                || state.contains_triggered(&cid)
2245                            {
2246                                log::debug!(
2247                                    "Skipping stale OrderStatusReport(Accepted) \
2248                                     for {cid} (order already terminal)"
2249                                );
2250                                continue;
2251                            }
2252
2253                            if !state.contains_accepted(&cid) {
2254                                state.insert_accepted(cid, order_report.venue_order_id);
2255                            }
2256                        }
2257                        OrderStatus::Triggered => {
2258                            if state.contains_filled(&cid) {
2259                                log::debug!(
2260                                    "Skipping stale OrderStatusReport(Triggered) \
2261                                     for {cid} (already filled)"
2262                                );
2263                                continue;
2264                            }
2265                            state.insert_triggered(cid);
2266                        }
2267                        OrderStatus::Filled => {
2268                            state.insert_filled(cid);
2269                            state.insert_terminal(cid);
2270                            state.remove_triggered(&cid);
2271                        }
2272                        OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected => {
2273                            state.insert_terminal(cid);
2274                            state.remove_triggered(&cid);
2275                            state.remove_filled(&cid);
2276                        }
2277                        _ => {}
2278                    }
2279                }
2280                emitter.send_order_status_report(order_report);
2281            }
2282            ExecutionReport::Fill(fill_report) => {
2283                if state.check_and_insert_trade(fill_report.trade_id) {
2284                    log::debug!(
2285                        "Skipping duplicate fill report: trade_id={}",
2286                        fill_report.trade_id
2287                    );
2288                    continue;
2289                }
2290
2291                if let Some(cid) = fill_report.client_order_id {
2292                    state.insert_filled(cid);
2293                    state.remove_triggered(&cid);
2294                }
2295                emitter.send_fill_report(fill_report);
2296            }
2297        }
2298    }
2299}
2300
2301fn emit_send_failed_submit(
2302    failure: &CommandFailure,
2303    state: &WsDispatchState,
2304    emitter: &ExecutionEventEmitter,
2305    clock: &AtomicTime,
2306    client_order_id: ClientOrderId,
2307) {
2308    let (CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason)) = failure else {
2309        return;
2310    };
2311    let Some(ident) = state.order_identity(client_order_id) else {
2312        return;
2313    };
2314
2315    state.remove_order_tracking(client_order_id);
2316    emitter.emit_order_rejected_event(
2317        ident.strategy_id,
2318        ident.instrument_id,
2319        client_order_id,
2320        reason,
2321        clock.get_time_ns(),
2322        false,
2323    );
2324}
2325
2326fn emit_send_failed_modify(
2327    failure: &CommandFailure,
2328    state: &WsDispatchState,
2329    emitter: &ExecutionEventEmitter,
2330    clock: &AtomicTime,
2331    client_order_id: ClientOrderId,
2332) {
2333    let (CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason)) = failure else {
2334        return;
2335    };
2336    let Some(ident) = state.order_identity(client_order_id) else {
2337        return;
2338    };
2339
2340    emitter.emit_order_modify_rejected_event(
2341        ident.strategy_id,
2342        ident.instrument_id,
2343        client_order_id,
2344        None,
2345        reason,
2346        clock.get_time_ns(),
2347    );
2348}
2349
2350fn format_order_response_reason(s_code: &str, s_msg: &str, sub_code: &str) -> String {
2351    match (s_msg.is_empty(), sub_code.is_empty(), s_code.is_empty()) {
2352        (false, true, _) => s_msg.to_string(),
2353        (false, false, _) => format!("{s_msg} (subCode={sub_code})"),
2354        (true, false, false) => format!("sCode={s_code} subCode={sub_code}"),
2355        (true, false, true) => format!("subCode={sub_code}"),
2356        (true, true, false) => format!("sCode={s_code}"),
2357        (true, true, true) => String::new(),
2358    }
2359}
2360
2361#[derive(Debug, Clone)]
2362pub struct AlgoCancelContext {
2363    pub client_order_id: ClientOrderId,
2364    pub instrument_id: InstrumentId,
2365    pub strategy_id: StrategyId,
2366    pub venue_order_id: Option<VenueOrderId>,
2367}
2368
2369// Contexts must correspond 1:1 with the requests that produced
2370// the responses (OKX preserves request order in batch responses).
2371pub fn emit_algo_cancel_rejections(
2372    responses: &[OKXCancelAlgoOrderResponse],
2373    contexts: &[AlgoCancelContext],
2374    emitter: &ExecutionEventEmitter,
2375    clock: &'static AtomicTime,
2376) {
2377    for (i, item) in responses.iter().enumerate() {
2378        let code = item.s_code.as_deref().unwrap_or(OKX_SUCCESS_CODE);
2379        if code == OKX_SUCCESS_CODE {
2380            continue;
2381        }
2382
2383        let msg = item.s_msg.as_deref().unwrap_or("");
2384
2385        if matches!(
2386            classify_okx_venue_code(code, msg),
2387            CommandFailure::Ambiguous(_) | CommandFailure::NotSent(_)
2388        ) {
2389            if let Some(ctx) = contexts.get(i) {
2390                log::warn!(
2391                    "Ambiguous algo cancel response for {}, awaiting reconciliation: \
2392                     algo_id={} sCode={code} sMsg={msg}",
2393                    ctx.client_order_id,
2394                    item.algo_id
2395                );
2396            } else {
2397                log::warn!(
2398                    "Ambiguous algo cancel response without context at index {i}: \
2399                     algo_id={} sCode={code} sMsg={msg}",
2400                    item.algo_id
2401                );
2402            }
2403            continue;
2404        }
2405
2406        if let Some(ctx) = contexts.get(i) {
2407            let ts = clock.get_time_ns();
2408            emitter.emit_order_cancel_rejected_event(
2409                ctx.strategy_id,
2410                ctx.instrument_id,
2411                ctx.client_order_id,
2412                ctx.venue_order_id,
2413                msg,
2414                ts,
2415            );
2416        } else {
2417            log::warn!(
2418                "Algo cancel rejected but no context at index {i}: \
2419                 algo_id={} sCode={code} sMsg={msg}",
2420                item.algo_id
2421            );
2422        }
2423    }
2424}
2425
2426pub fn emit_batch_cancel_failure(
2427    contexts: &[AlgoCancelContext],
2428    error: &str,
2429    _emitter: &ExecutionEventEmitter,
2430    _clock: &'static AtomicTime,
2431) {
2432    for ctx in contexts {
2433        log::warn!(
2434            "Ambiguous algo batch cancel failure for {}, awaiting reconciliation: {error}",
2435            ctx.client_order_id
2436        );
2437    }
2438}
2439
2440#[cfg(test)]
2441mod tests {
2442    use std::time::Duration;
2443
2444    use nautilus_common::messages::{ExecutionEvent, ExecutionReport as CommonExecutionReport};
2445    use nautilus_core::time::get_atomic_clock_realtime;
2446    use nautilus_model::{
2447        enums::{AccountType, OrderSide, OrderType, TimeInForce, TriggerType},
2448        identifiers::Symbol,
2449        instruments::CryptoPerpetual,
2450        types::Price,
2451    };
2452    use rstest::rstest;
2453
2454    use super::*;
2455    use crate::websocket::{error::OKXWsError, messages::OKXWsFrame};
2456
2457    fn load_algo_order_messages(fixture: &str) -> Vec<OKXAlgoOrderMsg> {
2458        let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2459            .join("test_data")
2460            .join(fixture);
2461        let content = std::fs::read_to_string(path).unwrap();
2462        let frame: OKXWsFrame = serde_json::from_str(&content).unwrap();
2463        let OKXWsFrame::Data { data, .. } = frame else {
2464            panic!("Expected algo order data frame");
2465        };
2466        serde_json::from_value(data).unwrap()
2467    }
2468
2469    fn load_regular_order_messages(fixture: &str) -> Vec<OKXOrderMsg> {
2470        let path = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"))
2471            .join("test_data")
2472            .join(fixture);
2473        let content = std::fs::read_to_string(path).unwrap();
2474        let frame: OKXWsFrame = serde_json::from_str(&content).unwrap();
2475        let OKXWsFrame::Data { data, .. } = frame else {
2476            panic!("Expected regular order data frame");
2477        };
2478        serde_json::from_value(data).unwrap()
2479    }
2480
2481    fn test_algo_context(client_order_id: ClientOrderId) -> OrderContext {
2482        OrderContext {
2483            identity: OrderIdentity {
2484                client_order_id,
2485                strategy_id: StrategyId::from("STRATEGY-001"),
2486                instrument_id: InstrumentId::from("BTC-USDT-SWAP.OKX"),
2487                order_side: OrderSide::Sell,
2488                order_type: OrderType::StopLimit,
2489            },
2490            quantity: Quantity::from("0.01"),
2491            price: Some(Price::from("102900")),
2492            trigger_price: Some(Price::from("95000")),
2493            trigger_type: Some(TriggerType::LastPrice),
2494            time_in_force: TimeInForce::Gtc,
2495            is_post_only: false,
2496            is_reduce_only: true,
2497            is_quote_quantity: false,
2498        }
2499    }
2500
2501    fn test_algo_instruments() -> AtomicMap<Ustr, InstrumentAny> {
2502        let instrument = CryptoPerpetual::builder()
2503            .instrument_id(InstrumentId::from("BTC-USDT-SWAP.OKX"))
2504            .raw_symbol(Symbol::from("BTC-USDT-SWAP"))
2505            .base_currency(Currency::BTC())
2506            .quote_currency(Currency::USDT())
2507            .settlement_currency(Currency::USDT())
2508            .is_inverse(false)
2509            .price_precision(2)
2510            .size_precision(8)
2511            .price_increment(Price::from("0.01"))
2512            .size_increment(Quantity::from("0.00000001"))
2513            .ts_event(UnixNanos::default())
2514            .ts_init(UnixNanos::default())
2515            .build()
2516            .unwrap();
2517        let instruments = AtomicMap::new();
2518        instruments.insert(
2519            Ustr::from("BTC-USDT-SWAP"),
2520            InstrumentAny::CryptoPerpetual(instrument),
2521        );
2522        instruments
2523    }
2524
2525    fn test_execution_emitter() -> (
2526        ExecutionEventEmitter,
2527        tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2528    ) {
2529        let clock = get_atomic_clock_realtime();
2530        let mut emitter = ExecutionEventEmitter::new(
2531            clock,
2532            TraderId::from("TRADER-001"),
2533            AccountId::from("OKX-001"),
2534            AccountType::Margin,
2535            None,
2536        );
2537        let (sender, receiver) = tokio::sync::mpsc::unbounded_channel();
2538        emitter.set_sender(sender);
2539        (emitter, receiver)
2540    }
2541
2542    fn dispatch_test_message(
2543        message: OKXWsMessage,
2544        emitter: &ExecutionEventEmitter,
2545        state: &WsDispatchState,
2546        instruments: &AtomicMap<Ustr, InstrumentAny>,
2547    ) {
2548        dispatch_ws_message(
2549            message,
2550            emitter,
2551            state,
2552            AccountId::from("OKX-001"),
2553            instruments,
2554            &mut AHashMap::new(),
2555            &mut AHashMap::new(),
2556            &mut AHashMap::new(),
2557            get_atomic_clock_realtime(),
2558        );
2559    }
2560
2561    fn drain_execution_events(
2562        receiver: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2563    ) -> Vec<ExecutionEvent> {
2564        let mut events = Vec::new();
2565        while let Ok(event) = receiver.try_recv() {
2566            events.push(event);
2567        }
2568        events
2569    }
2570
2571    #[rstest]
2572    fn tracked_algo_live_and_pause_emit_one_parent_acceptance() {
2573        let mut messages = load_algo_order_messages("ws_orders_algo.json");
2574        let live = messages.remove(0);
2575        let mut pause = live.clone();
2576        pause.state = OKXAlgoOrderStatus::Pause;
2577        pause.u_time += 1;
2578        let client_order_id = ClientOrderId::new(live.algo_cl_ord_id.as_str());
2579        let state = WsDispatchState::default();
2580        state.track_order_context(test_algo_context(client_order_id));
2581        let instruments = test_algo_instruments();
2582        let (emitter, mut receiver) = test_execution_emitter();
2583
2584        dispatch_test_message(
2585            OKXWsMessage::AlgoOrders(vec![live, pause]),
2586            &emitter,
2587            &state,
2588            &instruments,
2589        );
2590
2591        let events = drain_execution_events(&mut receiver);
2592        assert_eq!(events.len(), 1);
2593        match &events[0] {
2594            ExecutionEvent::Order(OrderEventAny::Accepted(accepted)) => {
2595                assert_eq!(accepted.client_order_id, client_order_id);
2596                assert_eq!(
2597                    accepted.venue_order_id,
2598                    VenueOrderId::new("706620792746729472")
2599                );
2600            }
2601            other => panic!("Expected tracked algo acceptance, was {other:?}"),
2602        }
2603    }
2604
2605    #[rstest]
2606    fn algo_update_routes_are_explicit_for_external_tracked_and_suppressed() {
2607        let mut message = load_algo_order_messages("ws_orders_algo.json").remove(0);
2608        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
2609        let state = WsDispatchState::default();
2610
2611        assert_eq!(
2612            route_algo_order_message(&message, &state),
2613            ExecutionUpdateRoute::External
2614        );
2615
2616        let context = test_algo_context(client_order_id);
2617        state.track_order_context(context);
2618        assert_eq!(
2619            route_algo_order_message(&message, &state),
2620            ExecutionUpdateRoute::Tracked(client_order_id, context)
2621        );
2622
2623        message.state = OKXAlgoOrderStatus::Unknown;
2624        assert_eq!(
2625            route_algo_order_message(&message, &state),
2626            ExecutionUpdateRoute::Suppressed
2627        );
2628    }
2629
2630    #[rstest]
2631    fn rest_parent_binding_routes_linked_child_before_parent_stream_update() {
2632        let mut child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2633        let client_order_id = ClientOrderId::new("STOP003BTCUSDT20250120");
2634        let parent_venue_order_id = VenueOrderId::new("706620792746729474");
2635        let child_venue_order_id = VenueOrderId::new(child.ord_id);
2636        child.algo_cl_ord_id = None;
2637        let state = WsDispatchState::default();
2638        state.track_order_context(test_algo_context(client_order_id));
2639        state.bind_algo_parent(client_order_id, parent_venue_order_id);
2640        let instruments = test_algo_instruments();
2641        let (emitter, mut receiver) = test_execution_emitter();
2642
2643        dispatch_test_message(
2644            OKXWsMessage::Orders(vec![child.clone()]),
2645            &emitter,
2646            &state,
2647            &instruments,
2648        );
2649
2650        let events = drain_execution_events(&mut receiver);
2651        assert_eq!(events.len(), 4);
2652        assert!(matches!(
2653            &events[0],
2654            ExecutionEvent::Order(OrderEventAny::Accepted(accepted))
2655                if accepted.client_order_id == client_order_id
2656                    && accepted.venue_order_id == parent_venue_order_id
2657        ));
2658        assert!(matches!(
2659            &events[1],
2660            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2661                if updated.client_order_id == client_order_id
2662                    && updated.venue_order_id == Some(child_venue_order_id)
2663        ));
2664        assert!(matches!(
2665            &events[2],
2666            ExecutionEvent::Order(OrderEventAny::Triggered(triggered))
2667                if triggered.client_order_id == client_order_id
2668                    && triggered.venue_order_id == Some(child_venue_order_id)
2669        ));
2670        assert!(matches!(
2671            &events[3],
2672            ExecutionEvent::Order(OrderEventAny::Filled(filled))
2673                if filled.client_order_id == client_order_id
2674                    && filled.venue_order_id == child_venue_order_id
2675        ));
2676
2677        dispatch_test_message(
2678            OKXWsMessage::Orders(vec![child]),
2679            &emitter,
2680            &state,
2681            &instruments,
2682        );
2683        assert!(drain_execution_events(&mut receiver).is_empty());
2684    }
2685
2686    #[rstest]
2687    #[tokio::test]
2688    async fn pre_binding_child_is_held_until_parent_binding() {
2689        let child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2690        let parent_venue_order_id = VenueOrderId::new(
2691            child
2692                .linked_algo_ord
2693                .as_ref()
2694                .expect("Expected linked algo order")
2695                .algo_id
2696                .as_str(),
2697        );
2698        let client_order_id = ClientOrderId::new("STOP003BTCUSDT20250120");
2699        let child_venue_order_id = VenueOrderId::new(child.ord_id);
2700        let state = WsDispatchState::default();
2701        state.track_order_context(test_algo_context(client_order_id));
2702        let instruments = test_algo_instruments();
2703        let (emitter, mut receiver) = test_execution_emitter();
2704
2705        dispatch_test_message(
2706            OKXWsMessage::Orders(vec![child.clone()]),
2707            &emitter,
2708            &state,
2709            &instruments,
2710        );
2711        assert!(drain_execution_events(&mut receiver).is_empty());
2712
2713        state.bind_algo_parent(client_order_id, parent_venue_order_id);
2714        tokio::time::timeout(
2715            Duration::from_millis(100),
2716            state.wait_for_linked_child_route(),
2717        )
2718        .await
2719        .expect("Expected parent binding to wake the private dispatch loop");
2720        dispatch_test_message(
2721            OKXWsMessage::Orders(Vec::new()),
2722            &emitter,
2723            &state,
2724            &instruments,
2725        );
2726
2727        let events = drain_execution_events(&mut receiver);
2728        assert_eq!(events.len(), 4);
2729        assert!(matches!(
2730            &events[0],
2731            ExecutionEvent::Order(OrderEventAny::Accepted(accepted))
2732                if accepted.client_order_id == client_order_id
2733                    && accepted.venue_order_id == parent_venue_order_id
2734        ));
2735        assert!(matches!(
2736            &events[1],
2737            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2738                if updated.client_order_id == client_order_id
2739                    && updated.venue_order_id == Some(child_venue_order_id)
2740        ));
2741        assert!(matches!(
2742            &events[2],
2743            ExecutionEvent::Order(OrderEventAny::Triggered(triggered))
2744                if triggered.client_order_id == client_order_id
2745                    && triggered.venue_order_id == Some(child_venue_order_id)
2746        ));
2747        assert!(matches!(
2748            &events[3],
2749            ExecutionEvent::Order(OrderEventAny::Filled(filled))
2750                if filled.client_order_id == client_order_id
2751                    && filled.venue_order_id == child_venue_order_id
2752        ));
2753
2754        dispatch_test_message(OKXWsMessage::Reconnected, &emitter, &state, &instruments);
2755        dispatch_test_message(
2756            OKXWsMessage::Orders(vec![child]),
2757            &emitter,
2758            &state,
2759            &instruments,
2760        );
2761        assert!(drain_execution_events(&mut receiver).is_empty());
2762    }
2763
2764    #[rstest]
2765    fn held_child_becomes_external_when_candidate_binds_another_parent() {
2766        let child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2767        let client_order_id = ClientOrderId::new("STOP003BTCUSDT20250120");
2768        let child_venue_order_id = VenueOrderId::new(child.ord_id);
2769        let state = WsDispatchState::default();
2770        state.track_order_context(test_algo_context(client_order_id));
2771        let instruments = test_algo_instruments();
2772        let (emitter, mut receiver) = test_execution_emitter();
2773
2774        dispatch_test_message(
2775            OKXWsMessage::Orders(vec![child]),
2776            &emitter,
2777            &state,
2778            &instruments,
2779        );
2780        assert!(drain_execution_events(&mut receiver).is_empty());
2781
2782        let authoritative_parent = VenueOrderId::new("706620792746729475");
2783        state.bind_algo_parent(client_order_id, authoritative_parent);
2784        dispatch_test_message(
2785            OKXWsMessage::Orders(Vec::new()),
2786            &emitter,
2787            &state,
2788            &instruments,
2789        );
2790
2791        let events = drain_execution_events(&mut receiver);
2792        assert_eq!(events.len(), 1);
2793        assert!(matches!(
2794            &events[0],
2795            ExecutionEvent::Report(CommonExecutionReport::Fill(report))
2796                if report.client_order_id
2797                    == Some(ClientOrderId::new("706620792746729474_0"))
2798                    && report.venue_order_id == child_venue_order_id
2799        ));
2800        assert_eq!(
2801            state.order_venue_binding(client_order_id),
2802            Some((authoritative_parent, false))
2803        );
2804    }
2805
2806    #[rstest]
2807    fn active_algo_bindings_are_not_evicted_with_replay_state() {
2808        let state = WsDispatchState::default();
2809        let first_client_order_id = ClientOrderId::new("STOP-ACTIVE-00000");
2810        let first_venue_order_id = VenueOrderId::new("700000000000000000");
2811
2812        for index in 0..=DEDUP_CAPACITY {
2813            let client_order_id = ClientOrderId::new(format!("STOP-ACTIVE-{index:05}").as_str());
2814            let venue_order_id = VenueOrderId::new(
2815                format!("{}", 700_000_000_000_000_000_u64 + index as u64).as_str(),
2816            );
2817            state.bind_algo_parent(client_order_id, venue_order_id);
2818        }
2819
2820        assert_eq!(
2821            state.order_venue_binding(first_client_order_id),
2822            Some((first_venue_order_id, false))
2823        );
2824    }
2825
2826    #[rstest]
2827    fn ambiguous_submit_failure_preserves_only_bound_context() {
2828        let client_order_id = ClientOrderId::new("STOP-AMBIGUOUS-001");
2829        let failure = CommandFailure::Ambiguous("request timed out".to_string());
2830        let unbound_state = WsDispatchState::default();
2831        unbound_state.track_order_context(test_algo_context(client_order_id));
2832
2833        unbound_state.resolve_algo_submit_failure(client_order_id, &failure);
2834        assert_eq!(unbound_state.order_identity(client_order_id), None);
2835
2836        let bound_state = WsDispatchState::default();
2837        let context = test_algo_context(client_order_id);
2838        let parent_venue_order_id = VenueOrderId::new("706620792746729476");
2839        bound_state.track_order_context(context);
2840        bound_state.bind_algo_parent(client_order_id, parent_venue_order_id);
2841
2842        bound_state.resolve_algo_submit_failure(client_order_id, &failure);
2843        assert_eq!(
2844            bound_state.order_identity(client_order_id),
2845            Some(context.identity)
2846        );
2847        assert_eq!(
2848            bound_state.order_venue_binding(client_order_id),
2849            Some((parent_venue_order_id, false))
2850        );
2851    }
2852
2853    #[rstest]
2854    #[case::effective(0, false)]
2855    #[case::effective_order_id_list(0, true)]
2856    #[case::partially_effective(1, false)]
2857    #[case::order_placed(2, false)]
2858    fn tracked_algo_trigger_states_bind_child_before_trigger(
2859        #[case] index: usize,
2860        #[case] use_order_id_list: bool,
2861    ) {
2862        let mut messages = if index == 2 {
2863            load_algo_order_messages("ws_orders_algo.json")
2864        } else {
2865            load_algo_order_messages("ws_orders_algo_states.json")
2866        };
2867        let mut message = if index == 2 {
2868            messages.remove(2)
2869        } else {
2870            messages.remove(index)
2871        };
2872        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
2873        let child_venue_order_id = VenueOrderId::new(message.ord_id.as_str());
2874        if use_order_id_list {
2875            message.ord_id_list = vec![message.ord_id.clone()];
2876            message.ord_id.clear();
2877        }
2878
2879        let state = WsDispatchState::default();
2880        state.track_order_context(test_algo_context(client_order_id));
2881        let instruments = test_algo_instruments();
2882        let (emitter, mut receiver) = test_execution_emitter();
2883
2884        dispatch_test_message(
2885            OKXWsMessage::AlgoOrders(vec![message]),
2886            &emitter,
2887            &state,
2888            &instruments,
2889        );
2890
2891        let events = drain_execution_events(&mut receiver);
2892        assert_eq!(events.len(), 3);
2893        assert!(matches!(
2894            &events[0],
2895            ExecutionEvent::Order(OrderEventAny::Accepted(_))
2896        ));
2897        assert!(matches!(
2898            &events[1],
2899            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2900                if updated.venue_order_id == Some(child_venue_order_id)
2901        ));
2902        assert!(matches!(
2903            &events[2],
2904            ExecutionEvent::Order(OrderEventAny::Triggered(triggered))
2905                if triggered.venue_order_id == Some(child_venue_order_id)
2906        ));
2907        assert_eq!(
2908            state.order_venue_binding(client_order_id),
2909            Some((child_venue_order_id, true))
2910        );
2911    }
2912
2913    #[rstest]
2914    fn tracked_algo_transition_uses_venue_actual_quantity() {
2915        let mut message = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
2916        message.sz.clear();
2917        message.close_fraction = "1".to_string();
2918        message.actual_sz = "0.025".to_string();
2919        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
2920        let state = WsDispatchState::default();
2921        state.track_order_context(test_algo_context(client_order_id));
2922        let instruments = test_algo_instruments();
2923        let (emitter, mut receiver) = test_execution_emitter();
2924
2925        dispatch_test_message(
2926            OKXWsMessage::AlgoOrders(vec![message]),
2927            &emitter,
2928            &state,
2929            &instruments,
2930        );
2931
2932        let events = drain_execution_events(&mut receiver);
2933        assert_eq!(events.len(), 3);
2934        assert!(matches!(
2935            &events[1],
2936            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2937                if updated.quantity == Quantity::from("0.025")
2938                    && updated.price.is_none()
2939        ));
2940        assert_eq!(
2941            state
2942                .order_contexts
2943                .get(&client_order_id)
2944                .map(|context| context.quantity),
2945            Some(Quantity::from("0.025"))
2946        );
2947    }
2948
2949    #[rstest]
2950    fn regular_child_refreshes_parent_transition_terms() {
2951        let mut parent = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
2952        parent.ord_px = "94950".to_string();
2953        let parent_replay = parent.clone();
2954        let client_order_id = ClientOrderId::new(parent.algo_cl_ord_id.as_str());
2955        let parent_venue_order_id = parent.algo_id.clone();
2956        let child_venue_order_id = parent.ord_id.clone();
2957        let state = WsDispatchState::default();
2958        state.track_order_context(test_algo_context(client_order_id));
2959        let instruments = test_algo_instruments();
2960        let (emitter, mut receiver) = test_execution_emitter();
2961
2962        dispatch_test_message(
2963            OKXWsMessage::AlgoOrders(vec![parent]),
2964            &emitter,
2965            &state,
2966            &instruments,
2967        );
2968        assert_eq!(drain_execution_events(&mut receiver).len(), 3);
2969
2970        let mut child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
2971        child.algo_id = Some(parent_venue_order_id.clone());
2972        child.algo_cl_ord_id = None;
2973        child.linked_algo_ord = Some(crate::websocket::messages::OKXLinkedAlgoOrd {
2974            algo_id: parent_venue_order_id,
2975        });
2976        child.ord_id = Ustr::from(child_venue_order_id.as_str());
2977        child.ord_type = OKXOrderType::Limit;
2978        child.state = OKXOrderStatus::Live;
2979        child.sz = "0.025".to_string();
2980        child.px = "94950".to_string();
2981        child.acc_fill_sz = Some("0".to_string());
2982        child.fill_sz.clear();
2983        child.fill_px.clear();
2984        child.trade_id.clear();
2985
2986        dispatch_test_message(
2987            OKXWsMessage::Orders(vec![child]),
2988            &emitter,
2989            &state,
2990            &instruments,
2991        );
2992
2993        let events = drain_execution_events(&mut receiver);
2994        assert_eq!(events.len(), 1);
2995        assert!(matches!(
2996            &events[0],
2997            ExecutionEvent::Order(OrderEventAny::Updated(updated))
2998                if updated.quantity == Quantity::from("0.025")
2999                    && updated.price == Some(Price::from("94950"))
3000        ));
3001        assert_eq!(
3002            state
3003                .order_contexts
3004                .get(&client_order_id)
3005                .map(|context| (context.quantity, context.price)),
3006            Some((Quantity::from("0.025"), Some(Price::from("94950"))))
3007        );
3008
3009        dispatch_test_message(
3010            OKXWsMessage::AlgoOrders(vec![parent_replay]),
3011            &emitter,
3012            &state,
3013            &instruments,
3014        );
3015        assert!(drain_execution_events(&mut receiver).is_empty());
3016        assert_eq!(
3017            state
3018                .order_contexts
3019                .get(&client_order_id)
3020                .map(|context| context.quantity),
3021            Some(Quantity::from("0.025"))
3022        );
3023    }
3024
3025    #[rstest]
3026    fn tracked_algo_holds_trigger_until_linked_child_arrives() {
3027        let mut parent = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
3028        parent.algo_id = "706620792746729474".to_string();
3029        parent.algo_cl_ord_id = "STOP003BTCUSDT20250120".to_string();
3030        parent.ord_id.clear();
3031        parent.ord_id_list.clear();
3032        let client_order_id = ClientOrderId::new(parent.algo_cl_ord_id.as_str());
3033        let state = WsDispatchState::default();
3034        state.track_order_context(test_algo_context(client_order_id));
3035        let instruments = test_algo_instruments();
3036        let (emitter, mut receiver) = test_execution_emitter();
3037
3038        dispatch_test_message(
3039            OKXWsMessage::AlgoOrders(vec![parent]),
3040            &emitter,
3041            &state,
3042            &instruments,
3043        );
3044
3045        let parent_events = drain_execution_events(&mut receiver);
3046        assert_eq!(parent_events.len(), 1);
3047        assert!(matches!(
3048            &parent_events[0],
3049            ExecutionEvent::Order(OrderEventAny::Accepted(_))
3050        ));
3051        assert!(!state.contains_triggered(&client_order_id));
3052
3053        let child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
3054        assert!(child.algo_cl_ord_id.is_none());
3055        assert_eq!(child.cl_ord_id, "706620792746729474_0");
3056        assert_eq!(
3057            child
3058                .linked_algo_ord
3059                .as_ref()
3060                .map(|linked| linked.algo_id.as_str()),
3061            Some("706620792746729474")
3062        );
3063        dispatch_test_message(
3064            OKXWsMessage::Orders(vec![child]),
3065            &emitter,
3066            &state,
3067            &instruments,
3068        );
3069
3070        let child_events = drain_execution_events(&mut receiver);
3071        assert_eq!(child_events.len(), 3);
3072        assert!(matches!(
3073            &child_events[0],
3074            ExecutionEvent::Order(OrderEventAny::Updated(_))
3075        ));
3076        assert!(matches!(
3077            &child_events[1],
3078            ExecutionEvent::Order(OrderEventAny::Triggered(_))
3079        ));
3080
3081        match &child_events[2] {
3082            ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
3083                assert_eq!(filled.client_order_id, client_order_id);
3084                assert_eq!(
3085                    filled.venue_order_id,
3086                    VenueOrderId::new("706620792746729999")
3087                );
3088                assert_eq!(filled.trade_id, TradeId::new("1518905530"));
3089            }
3090            other => panic!("Expected tracked child fill, was {other:?}"),
3091        }
3092
3093        assert!(state.contains_terminal(&client_order_id));
3094    }
3095
3096    #[rstest]
3097    fn tracked_algo_child_cancellation_routes_late_fill_report() {
3098        let parent = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
3099        let client_order_id = ClientOrderId::new(parent.algo_cl_ord_id.as_str());
3100        let parent_venue_order_id = parent.algo_id.clone();
3101        let child_venue_order_id = parent.ord_id.clone();
3102        let state = WsDispatchState::default();
3103        state.track_order_context(test_algo_context(client_order_id));
3104        let instruments = test_algo_instruments();
3105        let (emitter, mut receiver) = test_execution_emitter();
3106
3107        dispatch_test_message(
3108            OKXWsMessage::AlgoOrders(vec![parent]),
3109            &emitter,
3110            &state,
3111            &instruments,
3112        );
3113        assert_eq!(drain_execution_events(&mut receiver).len(), 3);
3114
3115        let mut child = load_regular_order_messages("ws_orders_trigger.json").remove(0);
3116        child.algo_id = Some(parent_venue_order_id.clone());
3117        child.linked_algo_ord = Some(crate::websocket::messages::OKXLinkedAlgoOrd {
3118            algo_id: parent_venue_order_id,
3119        });
3120        child.ord_id = Ustr::from(child_venue_order_id.as_str());
3121        let mut canceled = child.clone();
3122        canceled.state = OKXOrderStatus::Canceled;
3123        canceled.acc_fill_sz = Some("0".to_string());
3124        canceled.fill_sz.clear();
3125        canceled.fill_px.clear();
3126        canceled.trade_id.clear();
3127        dispatch_test_message(
3128            OKXWsMessage::Orders(vec![canceled]),
3129            &emitter,
3130            &state,
3131            &instruments,
3132        );
3133
3134        let events = drain_execution_events(&mut receiver);
3135        assert_eq!(events.len(), 1);
3136        assert!(matches!(
3137            &events[0],
3138            ExecutionEvent::Order(OrderEventAny::Canceled(canceled))
3139                if canceled.client_order_id == client_order_id
3140                    && canceled.venue_order_id
3141                        == Some(VenueOrderId::new(child_venue_order_id.as_str()))
3142        ));
3143        assert!(state.contains_terminal(&client_order_id));
3144
3145        dispatch_test_message(
3146            OKXWsMessage::Orders(vec![child.clone()]),
3147            &emitter,
3148            &state,
3149            &instruments,
3150        );
3151        let events = drain_execution_events(&mut receiver);
3152        assert_eq!(events.len(), 1);
3153        assert!(matches!(
3154            &events[0],
3155            ExecutionEvent::Report(CommonExecutionReport::Fill(report))
3156                if report.client_order_id == Some(client_order_id)
3157                    && report.venue_order_id
3158                        == VenueOrderId::new(child_venue_order_id.as_str())
3159                    && report.trade_id == TradeId::new("1518905530")
3160        ));
3161
3162        dispatch_test_message(
3163            OKXWsMessage::Orders(vec![child]),
3164            &emitter,
3165            &state,
3166            &instruments,
3167        );
3168        assert!(drain_execution_events(&mut receiver).is_empty());
3169    }
3170
3171    #[rstest]
3172    fn reconnect_replay_and_stale_parent_acceptance_are_suppressed() {
3173        let message = load_algo_order_messages("ws_orders_algo_states.json").remove(0);
3174        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3175        let state = WsDispatchState::default();
3176        state.track_order_context(test_algo_context(client_order_id));
3177        let instruments = test_algo_instruments();
3178        let (emitter, mut receiver) = test_execution_emitter();
3179
3180        dispatch_test_message(
3181            OKXWsMessage::AlgoOrders(vec![message.clone()]),
3182            &emitter,
3183            &state,
3184            &instruments,
3185        );
3186        assert_eq!(drain_execution_events(&mut receiver).len(), 3);
3187
3188        let mut stale_live = message.clone();
3189        stale_live.state = OKXAlgoOrderStatus::Live;
3190        stale_live.ord_id.clear();
3191        dispatch_test_message(OKXWsMessage::Reconnected, &emitter, &state, &instruments);
3192        dispatch_test_message(
3193            OKXWsMessage::AlgoOrders(vec![message, stale_live]),
3194            &emitter,
3195            &state,
3196            &instruments,
3197        );
3198
3199        assert!(drain_execution_events(&mut receiver).is_empty());
3200        assert_eq!(
3201            state
3202                .order_venue_binding(client_order_id)
3203                .map(|(venue_order_id, _)| venue_order_id),
3204            Some(VenueOrderId::new("706620792746730010"))
3205        );
3206    }
3207
3208    #[rstest]
3209    #[case::filled("ws_orders_algo.json", 4, 3)]
3210    #[case::partially_failed("ws_orders_algo_states.json", 4, 1)]
3211    fn tracked_algo_aggregate_state_never_enters_reconciliation(
3212        #[case] fixture: &str,
3213        #[case] index: usize,
3214        #[case] expected_order_events: usize,
3215    ) {
3216        let message = load_algo_order_messages(fixture).remove(index);
3217        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3218        let state = WsDispatchState::default();
3219        state.track_order_context(test_algo_context(client_order_id));
3220        let instruments = test_algo_instruments();
3221        let (emitter, mut receiver) = test_execution_emitter();
3222
3223        dispatch_test_message(
3224            OKXWsMessage::AlgoOrders(vec![message]),
3225            &emitter,
3226            &state,
3227            &instruments,
3228        );
3229
3230        let events = drain_execution_events(&mut receiver);
3231        assert_eq!(events.len(), expected_order_events);
3232        assert!(
3233            events
3234                .iter()
3235                .all(|event| matches!(event, ExecutionEvent::Order(_)))
3236        );
3237        assert!(state.order_contexts.contains_key(&client_order_id));
3238        assert!(!state.contains_terminal(&client_order_id));
3239    }
3240
3241    #[rstest]
3242    fn tracked_algo_cancel_is_terminal_and_stale_live_is_suppressed() {
3243        let mut message = load_algo_order_messages("ws_orders_algo.json").remove(3);
3244        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3245        let state = WsDispatchState::default();
3246        state.track_order_context(test_algo_context(client_order_id));
3247        let instruments = test_algo_instruments();
3248        let (emitter, mut receiver) = test_execution_emitter();
3249
3250        dispatch_test_message(
3251            OKXWsMessage::AlgoOrders(vec![message.clone()]),
3252            &emitter,
3253            &state,
3254            &instruments,
3255        );
3256        message.state = OKXAlgoOrderStatus::Live;
3257        message.algo_cl_ord_id.clear();
3258        message.cl_ord_id.clear();
3259        dispatch_test_message(
3260            OKXWsMessage::AlgoOrders(vec![message]),
3261            &emitter,
3262            &state,
3263            &instruments,
3264        );
3265
3266        let events = drain_execution_events(&mut receiver);
3267        assert_eq!(events.len(), 2);
3268        assert!(matches!(
3269            &events[0],
3270            ExecutionEvent::Order(OrderEventAny::Accepted(_))
3271        ));
3272        assert!(matches!(
3273            &events[1],
3274            ExecutionEvent::Order(OrderEventAny::Canceled(canceled))
3275                if canceled.client_order_id == client_order_id
3276        ));
3277        assert!(state.contains_terminal(&client_order_id));
3278        assert!(!state.order_contexts.contains_key(&client_order_id));
3279        assert_eq!(state.order_venue_binding(client_order_id), None);
3280    }
3281
3282    #[rstest]
3283    fn tracked_algo_failure_uses_typed_rejection() {
3284        let mut message = load_algo_order_messages("ws_orders_algo_states.json").remove(3);
3285        message.fail_code = "51008".to_string();
3286        let client_order_id = ClientOrderId::new(message.algo_cl_ord_id.as_str());
3287        let state = WsDispatchState::default();
3288        state.track_order_context(test_algo_context(client_order_id));
3289        let instruments = test_algo_instruments();
3290        let (emitter, mut receiver) = test_execution_emitter();
3291
3292        dispatch_test_message(
3293            OKXWsMessage::AlgoOrders(vec![message]),
3294            &emitter,
3295            &state,
3296            &instruments,
3297        );
3298
3299        let events = drain_execution_events(&mut receiver);
3300        assert_eq!(events.len(), 1);
3301        assert!(matches!(
3302            &events[0],
3303            ExecutionEvent::Order(OrderEventAny::Rejected(rejected))
3304                if rejected.client_order_id == client_order_id
3305                    && rejected.reason == Ustr::from("51008")
3306        ));
3307        assert!(state.contains_terminal(&client_order_id));
3308    }
3309
3310    #[rstest]
3311    #[case("51000", "Rejected", "", "Rejected")]
3312    #[case("51000", "Rejected", "51004", "Rejected (subCode=51004)")]
3313    #[case("51000", "", "51004", "sCode=51000 subCode=51004")]
3314    #[case("51000", "", "", "sCode=51000")]
3315    #[case("", "", "51004", "subCode=51004")]
3316    #[case("", "", "", "")]
3317    fn test_format_order_response_reason(
3318        #[case] s_code: &str,
3319        #[case] s_msg: &str,
3320        #[case] sub_code: &str,
3321        #[case] expected: &str,
3322    ) {
3323        assert_eq!(
3324            format_order_response_reason(s_code, s_msg, sub_code),
3325            expected
3326        );
3327    }
3328
3329    #[rstest]
3330    #[case::ambiguous(OKXWsError::SendFailed("connection reset".to_string()), true)]
3331    #[case::not_sent(OKXWsError::NoActiveClient, false)]
3332    fn send_failure_preserves_only_ambiguous_pending_orders(
3333        #[case] error: OKXWsError,
3334        #[case] expected_pending: bool,
3335    ) {
3336        let client_order_ids = [
3337            ClientOrderId::from("O-batch-pending-1"),
3338            ClientOrderId::from("O-batch-pending-2"),
3339        ];
3340        let state = WsDispatchState::default();
3341        let instrument_id = InstrumentId::from("ETH-USDT-SWAP.OKX");
3342        let strategy_id = StrategyId::from("STRATEGY-001");
3343
3344        for client_order_id in client_order_ids {
3345            state.pending_orders.insert(
3346                client_order_id.to_string(),
3347                PendingOrderInfo {
3348                    trader_id: TraderId::from("TRADER-001"),
3349                    strategy_id,
3350                    instrument_id,
3351                },
3352            );
3353            state.order_identities.insert(
3354                client_order_id,
3355                OrderIdentity {
3356                    client_order_id,
3357                    instrument_id,
3358                    strategy_id,
3359                    order_side: OrderSide::Buy,
3360                    order_type: OrderType::Limit,
3361                },
3362            );
3363        }
3364
3365        let clock = get_atomic_clock_realtime();
3366        let mut emitter = ExecutionEventEmitter::new(
3367            clock,
3368            TraderId::from("TRADER-001"),
3369            AccountId::from("OKX-001"),
3370            AccountType::Margin,
3371            None,
3372        );
3373        let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel();
3374        emitter.set_sender(sender);
3375        let instruments = AtomicMap::new();
3376        let mut fee_cache = AHashMap::new();
3377        let mut filled_qty_cache = AHashMap::new();
3378        let mut order_state_cache = AHashMap::new();
3379
3380        dispatch_ws_message(
3381            OKXWsMessage::SendFailed {
3382                request_id: "req-batch-send-failure".to_string(),
3383                client_order_ids: client_order_ids.to_vec(),
3384                op: Some(OKXWsOperation::BatchOrders),
3385                error,
3386            },
3387            &emitter,
3388            &state,
3389            AccountId::from("OKX-001"),
3390            &instruments,
3391            &mut fee_cache,
3392            &mut filled_qty_cache,
3393            &mut order_state_cache,
3394            clock,
3395        );
3396
3397        for client_order_id in client_order_ids {
3398            assert_eq!(
3399                state.pending_orders.contains_key(client_order_id.as_str()),
3400                expected_pending,
3401                "pending state mismatch for {client_order_id}"
3402            );
3403        }
3404    }
3405}