1use 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
80const 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#[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 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 #[must_use]
399 pub fn contains_accepted(&self, cid: &ClientOrderId) -> bool {
400 self.accepted_venue_order_ids.lock().contains_key(cid)
401 }
402
403 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 #[must_use]
416 pub fn contains_triggered(&self, cid: &ClientOrderId) -> bool {
417 self.triggered_orders.contains(cid)
418 }
419
420 pub fn insert_triggered(&self, cid: ClientOrderId) {
422 let _ = self.triggered_orders.insert(cid);
423 }
424
425 #[must_use]
427 pub fn contains_filled(&self, cid: &ClientOrderId) -> bool {
428 self.filled_orders.contains(cid)
429 }
430
431 pub fn insert_filled(&self, cid: ClientOrderId) {
433 let _ = self.filled_orders.insert(cid);
434 }
435
436 #[must_use]
438 pub fn contains_terminal(&self, cid: &ClientOrderId) -> bool {
439 self.terminal_orders.contains(cid)
440 }
441
442 pub fn insert_terminal(&self, cid: ClientOrderId) {
444 let _ = self.terminal_orders.insert(cid);
445 }
446
447 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#[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#[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 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, );
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#[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 }
1951}
1952
1953#[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
2037fn 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#[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
2160fn 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
2226pub 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 #[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
2369pub 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}