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