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