1use std::sync::{
24 Arc,
25 atomic::{AtomicBool, Ordering},
26};
27
28use ahash::AHashMap;
29use anyhow::Context;
30use arc_swap::ArcSwapOption;
31use dashmap::{DashMap, DashSet};
32use nautilus_core::{UUID4, UnixNanos, time::AtomicTime};
33use nautilus_live::ExecutionEventEmitter;
34use nautilus_model::{
35 enums::{LiquiditySide, OrderSide, OrderType},
36 events::{
37 OrderAccepted, OrderCanceled, OrderEventAny, OrderFilled, OrderTriggered, OrderUpdated,
38 },
39 identifiers::{
40 AccountId, ClientOrderId, InstrumentId, PositionId, StrategyId, TradeId, VenueOrderId,
41 },
42 instruments::{Instrument, InstrumentAny},
43 orders::TRIGGERABLE_ORDER_TYPES,
44 types::{Money, Price, Quantity},
45};
46use rust_decimal::Decimal;
47use ustr::Ustr;
48
49use super::{
50 messages::{
51 BybitWsAccountExecution, BybitWsAccountExecutionFast, BybitWsAccountOrder, BybitWsMessage,
52 },
53 parse::{
54 parse_millis_i64, parse_ws_account_state, parse_ws_fill_report_fast,
55 parse_ws_position_status_report,
56 },
57};
58use crate::{
59 common::{
60 enums::{BybitExecType, BybitOrderSide, BybitOrderStatus, BybitProductType},
61 parse::{
62 bybit_rejection_due_post_only, get_currency, make_bybit_symbol, parse_millis_timestamp,
63 parse_price_with_precision, parse_quantity_with_precision,
64 },
65 },
66 http::error::is_bybit_ambiguous_order_error_code,
67 repay::RepayRequest,
68};
69
70const DEDUP_CAPACITY: usize = 10_000;
71
72const BYBIT_OP_ORDER_CREATE: &str = "order.create";
73const BYBIT_OP_ORDER_AMEND: &str = "order.amend";
74const BYBIT_OP_ORDER_CANCEL: &str = "order.cancel";
75const BYBIT_OP_ORDER_CREATE_BATCH: &str = "order.create-batch";
76const BYBIT_OP_ORDER_AMEND_BATCH: &str = "order.amend-batch";
77const BYBIT_OP_ORDER_CANCEL_BATCH: &str = "order.cancel-batch";
78
79#[derive(Debug, Clone)]
82pub struct OrderIdentity {
83 pub instrument_id: InstrumentId,
84 pub strategy_id: StrategyId,
85 pub order_side: OrderSide,
86 pub order_type: OrderType,
87 pub venue_position_id: Option<PositionId>,
88}
89
90#[derive(Debug, Clone, Copy)]
92pub enum PendingOperation {
93 Place,
94 Cancel,
95 Amend,
96}
97
98pub type PendingRequestData = (
101 Vec<ClientOrderId>,
102 Vec<Option<VenueOrderId>>,
103 PendingOperation,
104);
105
106#[derive(Debug, Clone)]
110pub struct OrderStateSnapshot {
111 pub quantity: Quantity,
112 pub price: Option<Price>,
113 pub trigger_price: Option<Price>,
114}
115
116#[derive(Clone, Copy, Debug)]
117struct SpotRepayFill {
118 quantity: Quantity,
119 base_fee: Decimal,
120}
121
122#[derive(Debug)]
123pub struct WsDispatchState {
124 pub order_identities: DashMap<ClientOrderId, OrderIdentity>,
125 pub pending_requests: DashMap<String, PendingRequestData>,
126 pub order_snapshots: DashMap<ClientOrderId, OrderStateSnapshot>,
127 pub emitted_accepted: DashSet<ClientOrderId>,
128 pub triggered_orders: DashSet<ClientOrderId>,
129 pub filled_orders: DashSet<ClientOrderId>,
130 spot_repay_fills: DashMap<ClientOrderId, SpotRepayFill>,
131 repay_tx: ArcSwapOption<tokio::sync::mpsc::UnboundedSender<RepayRequest>>,
132 clearing: AtomicBool,
133}
134
135impl Default for WsDispatchState {
136 fn default() -> Self {
137 Self {
138 order_identities: DashMap::new(),
139 pending_requests: DashMap::new(),
140 order_snapshots: DashMap::new(),
141 emitted_accepted: DashSet::default(),
142 triggered_orders: DashSet::default(),
143 filled_orders: DashSet::default(),
144 spot_repay_fills: DashMap::new(),
145 repay_tx: ArcSwapOption::empty(),
146 clearing: AtomicBool::new(false),
147 }
148 }
149}
150
151impl WsDispatchState {
152 fn evict_if_full(&self, set: &DashSet<ClientOrderId>) {
153 if set.len() >= DEDUP_CAPACITY
154 && self
155 .clearing
156 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
157 .is_ok()
158 {
159 set.clear();
160 self.clearing.store(false, Ordering::Release);
161 }
162 }
163
164 fn insert_accepted(&self, cid: ClientOrderId) {
165 self.evict_if_full(&self.emitted_accepted);
166 self.emitted_accepted.insert(cid);
167 }
168
169 fn insert_filled(&self, cid: ClientOrderId) {
170 self.evict_if_full(&self.filled_orders);
171 self.filled_orders.insert(cid);
172 }
173
174 fn insert_triggered(&self, cid: ClientOrderId) {
175 self.evict_if_full(&self.triggered_orders);
176 self.triggered_orders.insert(cid);
177 }
178
179 pub(crate) fn set_repay_sender(&self, tx: tokio::sync::mpsc::UnboundedSender<RepayRequest>) {
180 self.repay_tx.store(Some(Arc::new(tx)));
181 }
182
183 pub(crate) fn clear_repay_sender(&self) {
184 self.repay_tx.store(None);
185 }
186
187 fn enqueue_repay(&self, req: RepayRequest) {
188 if let Some(tx) = self.repay_tx.load_full()
189 && let Err(e) = tx.send(req)
190 {
191 log::warn!("Failed to enqueue spot borrow repayment: {e}");
192 }
193 }
194}
195
196pub fn dispatch_ws_message(
203 message: &BybitWsMessage,
204 emitter: &ExecutionEventEmitter,
205 state: &WsDispatchState,
206 account_id: AccountId,
207 instruments: &AHashMap<Ustr, InstrumentAny>,
208 clock: &AtomicTime,
209) {
210 match message {
211 BybitWsMessage::AccountOrder(msg) => {
212 let ts_init = clock.get_time_ns();
213
214 for order in &msg.data {
215 let symbol = make_bybit_symbol(order.symbol, order.category);
216 let Some(instrument) = instruments.get(&symbol) else {
217 log::warn!("No instrument for order update: {symbol}");
218 continue;
219 };
220 dispatch_order_update(order, instrument, emitter, state, account_id, ts_init);
221 }
222 }
223 BybitWsMessage::AccountExecution(msg) => {
224 let ts_init = clock.get_time_ns();
225
226 for exec in &msg.data {
227 let symbol = make_bybit_symbol(exec.symbol, exec.category);
228 let Some(instrument) = instruments.get(&symbol) else {
229 log::warn!("No instrument for execution update: {symbol}");
230 continue;
231 };
232 dispatch_execution_fill(exec, instrument, emitter, state, account_id, ts_init);
233 }
234 }
235 BybitWsMessage::AccountExecutionFast(msg) => {
236 let ts_init = clock.get_time_ns();
237
238 for exec in &msg.data {
239 let symbol = make_bybit_symbol(exec.symbol, exec.category);
240 let Some(instrument) = instruments.get(&symbol) else {
241 log::warn!("No instrument for fast-execution update: {symbol}");
242 continue;
243 };
244 dispatch_execution_fill_fast(exec, instrument, emitter, state, account_id, ts_init);
245 }
246 }
247 BybitWsMessage::AccountWallet(msg) => {
248 let ts_init = clock.get_time_ns();
249 let ts_event = parse_millis_i64(msg.creation_time, "wallet.creation_time")
250 .unwrap_or_else(|e| {
251 log::warn!("Failed to parse wallet creation_time, using ts_init: {e}");
252 ts_init
253 });
254
255 for wallet in &msg.data {
256 match parse_ws_account_state(wallet, account_id, ts_event, ts_init) {
257 Ok(state) => emitter.send_account_state(state),
258 Err(e) => log::error!("Failed to parse account state: {e}"),
259 }
260 }
261 }
262 BybitWsMessage::AccountPosition(msg) => {
263 let ts_init = clock.get_time_ns();
264
265 for position in &msg.data {
266 let symbol = make_bybit_symbol(position.symbol, position.category);
267 let Some(instrument) = instruments.get(&symbol) else {
268 log::warn!("No instrument for position update: {symbol}");
269 continue;
270 };
271
272 match parse_ws_position_status_report(position, account_id, instrument, ts_init) {
273 Ok(report) => emitter.send_position_report(report),
274 Err(e) => log::error!("Failed to parse position status report: {e}"),
275 }
276 }
277 }
278 BybitWsMessage::OrderResponse(resp) => {
279 let ts_init = clock.get_time_ns();
280 dispatch_order_response(resp, emitter, state, ts_init);
281 }
282 BybitWsMessage::Error(e) => {
283 log::warn!("WebSocket error: code={} message={}", e.code, e.message);
284 }
285 BybitWsMessage::Reconnected => {
286 log::info!("WebSocket reconnected");
287 }
288 BybitWsMessage::Auth(_)
289 | BybitWsMessage::Orderbook(_)
290 | BybitWsMessage::Trade(_)
291 | BybitWsMessage::Kline(_)
292 | BybitWsMessage::TickerLinear(_)
293 | BybitWsMessage::TickerOption(_) => {}
294 }
295}
296
297fn dispatch_order_update(
303 order: &BybitWsAccountOrder,
304 instrument: &InstrumentAny,
305 emitter: &ExecutionEventEmitter,
306 state: &WsDispatchState,
307 account_id: AccountId,
308 ts_init: UnixNanos,
309) {
310 let client_order_id = if order.order_link_id.is_empty() {
311 None
312 } else {
313 Some(ClientOrderId::new(order.order_link_id))
314 };
315
316 let identity = client_order_id
317 .as_ref()
318 .and_then(|cid| state.order_identities.get(cid).map(|r| r.clone()));
319
320 if let (Some(client_order_id), Some(identity)) = (client_order_id, identity) {
321 let venue_order_id = VenueOrderId::new(order.order_id);
322
323 match order.order_status {
324 BybitOrderStatus::Created | BybitOrderStatus::New | BybitOrderStatus::Untriggered => {
325 let snapshot = parse_order_snapshot(order, instrument);
326
327 if state.emitted_accepted.contains(&client_order_id)
328 || state.filled_orders.contains(&client_order_id)
329 || state.triggered_orders.contains(&client_order_id)
330 {
331 if let Some(snapshot) = snapshot
332 && is_snapshot_updated(&snapshot, &client_order_id, state)
333 {
334 let updated = OrderUpdated::new(
335 emitter.trader_id(),
336 identity.strategy_id,
337 identity.instrument_id,
338 client_order_id,
339 snapshot.quantity,
340 UUID4::new(),
341 ts_init,
342 ts_init,
343 false,
344 Some(venue_order_id),
345 Some(account_id),
346 snapshot.price,
347 snapshot.trigger_price,
348 None,
349 false,
350 );
351 state.order_snapshots.insert(client_order_id, snapshot);
352 emitter.send_order_event(OrderEventAny::Updated(updated));
353 return;
354 }
355 log::debug!("Skipping duplicate Accepted for {client_order_id}");
356 return;
357 }
358
359 state.insert_accepted(client_order_id);
360
361 let venue_differs_from_submitted = snapshot
364 .as_ref()
365 .is_some_and(|s| is_snapshot_updated(s, &client_order_id, state));
366
367 if let Some(snapshot) = snapshot.as_ref() {
368 state
369 .order_snapshots
370 .insert(client_order_id, snapshot.clone());
371 }
372
373 let accepted = OrderAccepted::new(
374 emitter.trader_id(),
375 identity.strategy_id,
376 identity.instrument_id,
377 client_order_id,
378 venue_order_id,
379 account_id,
380 UUID4::new(),
381 ts_init,
382 ts_init,
383 false,
384 );
385 emitter.send_order_event(OrderEventAny::Accepted(accepted));
386
387 if venue_differs_from_submitted && let Some(snapshot) = snapshot {
388 let updated = OrderUpdated::new(
389 emitter.trader_id(),
390 identity.strategy_id,
391 identity.instrument_id,
392 client_order_id,
393 snapshot.quantity,
394 UUID4::new(),
395 ts_init,
396 ts_init,
397 false,
398 Some(venue_order_id),
399 Some(account_id),
400 snapshot.price,
401 snapshot.trigger_price,
402 None,
403 false,
404 );
405 emitter.send_order_event(OrderEventAny::Updated(updated));
406 }
407 }
408 BybitOrderStatus::Triggered => {
409 if state.filled_orders.contains(&client_order_id) {
410 log::debug!("Skipping stale Triggered for {client_order_id} (already filled)");
411 return;
412 }
413
414 if !TRIGGERABLE_ORDER_TYPES.contains(&identity.order_type) {
415 log::debug!(
416 "Skipping OrderTriggered for {} order {client_order_id}: market-style stops have no TRIGGERED state",
417 identity.order_type,
418 );
419 return;
420 }
421
422 ensure_accepted_emitted(
423 client_order_id,
424 account_id,
425 venue_order_id,
426 &identity,
427 emitter,
428 state,
429 ts_init,
430 );
431 state.insert_triggered(client_order_id);
432 let triggered = OrderTriggered::new(
433 emitter.trader_id(),
434 identity.strategy_id,
435 identity.instrument_id,
436 client_order_id,
437 UUID4::new(),
438 ts_init,
439 ts_init,
440 false,
441 Some(venue_order_id),
442 Some(account_id),
443 );
444 emitter.send_order_event(OrderEventAny::Triggered(triggered));
445 }
446 BybitOrderStatus::Rejected => {
447 let filled_qty = parse_quantity_with_precision(
448 &order.cum_exec_qty,
449 instrument.size_precision(),
450 "order.cumExecQty",
451 )
452 .unwrap_or_default();
453
454 if filled_qty.is_positive() {
455 ensure_accepted_emitted(
457 client_order_id,
458 account_id,
459 venue_order_id,
460 &identity,
461 emitter,
462 state,
463 ts_init,
464 );
465 let canceled = OrderCanceled::new(
466 emitter.trader_id(),
467 identity.strategy_id,
468 identity.instrument_id,
469 client_order_id,
470 UUID4::new(),
471 ts_init,
472 ts_init,
473 false,
474 Some(venue_order_id),
475 Some(account_id),
476 None,
477 );
478 cleanup_terminal(client_order_id, state);
479 emitter.send_order_event(OrderEventAny::Canceled(canceled));
480 } else {
481 let reason = if order.reject_reason.is_empty() {
482 Ustr::from("Order rejected by venue")
483 } else {
484 order.reject_reason
485 };
486 state.order_identities.remove(&client_order_id);
487 state.order_snapshots.remove(&client_order_id);
488 emitter.emit_order_rejected_event(
489 identity.strategy_id,
490 identity.instrument_id,
491 client_order_id,
492 reason.as_str(),
493 ts_init,
494 bybit_rejection_due_post_only(reason.as_str()),
495 );
496 }
497 }
498 BybitOrderStatus::PartiallyFilled => {
499 ensure_accepted_emitted(
502 client_order_id,
503 account_id,
504 venue_order_id,
505 &identity,
506 emitter,
507 state,
508 ts_init,
509 );
510
511 if let Some(snapshot) = parse_order_snapshot(order, instrument)
515 && is_snapshot_updated(&snapshot, &client_order_id, state)
516 {
517 let updated = OrderUpdated::new(
518 emitter.trader_id(),
519 identity.strategy_id,
520 identity.instrument_id,
521 client_order_id,
522 snapshot.quantity,
523 UUID4::new(),
524 ts_init,
525 ts_init,
526 false,
527 Some(venue_order_id),
528 Some(account_id),
529 snapshot.price,
530 snapshot.trigger_price,
531 None,
532 false,
533 );
534 state.order_snapshots.insert(client_order_id, snapshot);
535 emitter.send_order_event(OrderEventAny::Updated(updated));
536 }
537 }
538 BybitOrderStatus::Filled => {
539 ensure_accepted_emitted(
542 client_order_id,
543 account_id,
544 venue_order_id,
545 &identity,
546 emitter,
547 state,
548 ts_init,
549 );
550
551 if let Some(snapshot) = parse_order_snapshot(order, instrument)
554 && is_snapshot_updated(&snapshot, &client_order_id, state)
555 {
556 let updated = OrderUpdated::new(
557 emitter.trader_id(),
558 identity.strategy_id,
559 identity.instrument_id,
560 client_order_id,
561 snapshot.quantity,
562 UUID4::new(),
563 ts_init,
564 ts_init,
565 false,
566 Some(venue_order_id),
567 Some(account_id),
568 snapshot.price,
569 snapshot.trigger_price,
570 None,
571 false,
572 );
573 state.order_snapshots.insert(client_order_id, snapshot);
574 emitter.send_order_event(OrderEventAny::Updated(updated));
575 }
576 }
580 BybitOrderStatus::Canceled
581 | BybitOrderStatus::PartiallyFilledCanceled
582 | BybitOrderStatus::Deactivated => {
583 let filled_qty = parse_quantity_with_precision(
584 &order.cum_exec_qty,
585 instrument.size_precision(),
586 "order.cumExecQty",
587 )
588 .unwrap_or_default();
589
590 if filled_qty.is_zero()
594 && bybit_rejection_due_post_only(order.reject_reason.as_str())
595 {
596 cleanup_terminal(client_order_id, state);
597 emitter.emit_order_rejected_event(
598 identity.strategy_id,
599 identity.instrument_id,
600 client_order_id,
601 order.reject_reason.as_str(),
602 ts_init,
603 true,
604 );
605 } else {
606 ensure_accepted_emitted(
607 client_order_id,
608 account_id,
609 venue_order_id,
610 &identity,
611 emitter,
612 state,
613 ts_init,
614 );
615 let canceled = OrderCanceled::new(
616 emitter.trader_id(),
617 identity.strategy_id,
618 identity.instrument_id,
619 client_order_id,
620 UUID4::new(),
621 ts_init,
622 ts_init,
623 false,
624 Some(venue_order_id),
625 Some(account_id),
626 None,
627 );
628 cleanup_terminal(client_order_id, state);
629 emitter.send_order_event(OrderEventAny::Canceled(canceled));
630 }
631 }
632 }
633 } else {
634 match super::parse::parse_ws_order_status_report(order, instrument, account_id, ts_init) {
636 Ok(report) => emitter.send_order_status_report(report),
637 Err(e) => log::error!("Failed to parse order status report: {e}"),
638 }
639 }
640}
641
642fn dispatch_execution_fill(
647 exec: &BybitWsAccountExecution,
648 instrument: &InstrumentAny,
649 emitter: &ExecutionEventEmitter,
650 state: &WsDispatchState,
651 account_id: AccountId,
652 ts_init: UnixNanos,
653) {
654 if exec.exec_type == BybitExecType::Funding {
655 log::debug!(
656 "Skipping funding execution: symbol={}, order_id={}, exec_id={}",
657 exec.symbol,
658 exec.order_id,
659 exec.exec_id,
660 );
661 return;
662 }
663
664 if exec.exec_type.is_exchange_generated() {
665 log::warn!(
666 "Exchange-generated execution: exec_type={:?}, symbol={}, order_id={}, order_link_id={}, side={:?}, qty={}, price={}",
667 exec.exec_type,
668 exec.symbol,
669 exec.order_id,
670 exec.order_link_id,
671 exec.side,
672 exec.exec_qty,
673 exec.exec_price,
674 );
675 }
676
677 let client_order_id = if exec.order_link_id.is_empty() {
678 None
679 } else {
680 Some(ClientOrderId::new(exec.order_link_id))
681 };
682
683 let identity = client_order_id
684 .as_ref()
685 .and_then(|cid| state.order_identities.get(cid).map(|r| r.clone()));
686
687 if let (Some(client_order_id), Some(identity)) = (client_order_id, identity) {
688 let venue_order_id = VenueOrderId::new(exec.order_id);
689
690 ensure_accepted_emitted(
691 client_order_id,
692 account_id,
693 venue_order_id,
694 &identity,
695 emitter,
696 state,
697 ts_init,
698 );
699
700 match parse_order_filled(exec, instrument, &identity, emitter, account_id, ts_init) {
701 Ok(filled) => {
702 let is_spot_buy =
703 exec.category == BybitProductType::Spot && exec.side == BybitOrderSide::Buy;
704
705 if is_spot_buy {
706 record_spot_repay_fill(client_order_id, &filled, instrument, state);
707 }
708
709 state.insert_filled(client_order_id);
710 state.triggered_orders.remove(&client_order_id);
711 emitter.send_order_event(OrderEventAny::Filled(filled));
712
713 if exec.leaves_qty == "0" {
714 if is_spot_buy {
715 enqueue_spot_repay(client_order_id, instrument, state);
716 }
717 cleanup_terminal(client_order_id, state);
718 }
719 }
720 Err(e) => log::error!("Failed to parse OrderFilled for {client_order_id}: {e}"),
721 }
722 } else {
723 match super::parse::parse_ws_fill_report(exec, account_id, instrument, ts_init) {
725 Ok(report) => emitter.send_fill_report(report),
726 Err(e) => log::error!("Failed to parse fill report: {e}"),
727 }
728 }
729}
730
731fn dispatch_execution_fill_fast(
738 exec: &BybitWsAccountExecutionFast,
739 instrument: &InstrumentAny,
740 emitter: &ExecutionEventEmitter,
741 state: &WsDispatchState,
742 account_id: AccountId,
743 ts_init: UnixNanos,
744) {
745 let client_order_id = if exec.order_link_id.is_empty() {
746 None
747 } else {
748 Some(ClientOrderId::new(exec.order_link_id))
749 };
750
751 let mut venue_position_id = None;
752
753 if let Some(cid) = client_order_id.as_ref()
754 && let Some(identity) = state.order_identities.get(cid).map(|r| r.clone())
755 {
756 venue_position_id = identity.venue_position_id;
757 let venue_order_id = VenueOrderId::new(exec.order_id);
758 ensure_accepted_emitted(
759 *cid,
760 account_id,
761 venue_order_id,
762 &identity,
763 emitter,
764 state,
765 ts_init,
766 );
767 }
768
769 match parse_ws_fill_report_fast(exec, account_id, instrument, venue_position_id, ts_init) {
770 Ok(report) => emitter.send_fill_report(report),
771 Err(e) => log::error!("Failed to parse fast fill report: {e}"),
772 }
773}
774
775fn parse_order_filled(
777 exec: &BybitWsAccountExecution,
778 instrument: &InstrumentAny,
779 identity: &OrderIdentity,
780 emitter: &ExecutionEventEmitter,
781 account_id: AccountId,
782 ts_init: UnixNanos,
783) -> anyhow::Result<OrderFilled> {
784 let client_order_id = ClientOrderId::new(exec.order_link_id);
785 let venue_order_id = VenueOrderId::new(exec.order_id);
786 let trade_id =
787 TradeId::new_checked(exec.exec_id.as_str()).context("invalid execId in Bybit execution")?;
788
789 let last_qty = parse_quantity_with_precision(
790 &exec.exec_qty,
791 instrument.size_precision(),
792 "execution.execQty",
793 )?;
794 let last_px = parse_price_with_precision(
795 &exec.exec_price,
796 instrument.price_precision(),
797 "execution.execPrice",
798 )?;
799
800 let liquidity_side = if exec.is_maker {
801 LiquiditySide::Maker
802 } else {
803 LiquiditySide::Taker
804 };
805
806 let fee_decimal: Decimal = exec
807 .exec_fee
808 .parse()
809 .with_context(|| format!("failed to parse execFee='{}'", exec.exec_fee))?;
810 let commission_currency = get_currency(&exec.fee_currency);
811 let commission = Money::from_decimal(fee_decimal, commission_currency).with_context(|| {
812 format!(
813 "failed to create commission from execFee='{}'",
814 exec.exec_fee
815 )
816 })?;
817
818 let ts_event = parse_millis_timestamp(&exec.exec_time, "execution.execTime")?;
819
820 Ok(OrderFilled::new(
821 emitter.trader_id(),
822 identity.strategy_id,
823 identity.instrument_id,
824 client_order_id,
825 venue_order_id,
826 account_id,
827 trade_id,
828 identity.order_side,
829 identity.order_type,
830 last_qty,
831 last_px,
832 commission_currency,
833 liquidity_side,
834 UUID4::new(),
835 ts_event,
836 ts_init,
837 false,
838 identity.venue_position_id,
839 Some(commission),
840 None,
841 ))
842}
843
844fn dispatch_order_response(
846 resp: &super::messages::BybitWsOrderResponse,
847 emitter: &ExecutionEventEmitter,
848 state: &WsDispatchState,
849 ts_init: UnixNanos,
850) {
851 if resp.ret_code == 0 {
852 let pending = resp
854 .req_id
855 .as_ref()
856 .and_then(|rid| state.pending_requests.remove(rid))
857 .map(|(_, v)| v);
858
859 if let Some((cids, voids, pending_op)) = pending {
860 let batch_errors = resp.extract_batch_errors();
861 let data_array = resp.data.as_array();
862
863 for (idx, error) in batch_errors.iter().enumerate() {
864 if error.code == 0 {
865 continue;
866 }
867
868 if is_bybit_ambiguous_order_error_code(error.code) {
869 log::warn!(
870 "Ambiguous batch order item failure at index {idx}: code={}, msg={}; awaiting reconciliation",
871 error.code,
872 error.msg,
873 );
874 continue;
875 }
876
877 let cid = data_array
879 .and_then(|arr| arr.get(idx))
880 .and_then(extract_order_link_id_from_data)
881 .or_else(|| cids.get(idx).copied());
882
883 let Some(cid) = cid else {
884 log::warn!(
885 "Batch error at index {idx} without correlation: code={}, msg={}",
886 error.code,
887 error.msg,
888 );
889 continue;
890 };
891
892 let Some(identity) = state.order_identities.get(&cid).map(|r| r.clone()) else {
893 log::warn!(
894 "Batch error for untracked order: client_order_id={cid}, msg={}",
895 error.msg,
896 );
897 continue;
898 };
899
900 let stored_void = voids.get(idx).and_then(|v| *v);
901
902 emit_rejection_for_op(
903 &pending_op,
904 cid,
905 &identity,
906 stored_void,
907 &error.msg,
908 emitter,
909 state,
910 ts_init,
911 );
912 }
913 }
914 return;
915 }
916
917 let pending = resp
919 .req_id
920 .as_ref()
921 .and_then(|rid| state.pending_requests.remove(rid))
922 .map(|(_, v)| v);
923
924 let is_batch_response = is_batch_order_op(resp.op.as_str())
925 || pending.as_ref().is_some_and(|(cids, _, _)| cids.len() > 1);
926
927 if is_batch_response {
928 if is_bybit_ambiguous_order_error_code(resp.ret_code) {
929 let order_count = pending.as_ref().map_or(0, |(cids, _, _)| cids.len());
930 log::warn!(
931 "Ambiguous batch order response failure for {order_count} orders: op={}, ret_code={}, ret_msg={}; awaiting reconciliation",
932 resp.op,
933 resp.ret_code,
934 resp.ret_msg,
935 );
936 return;
937 }
938
939 let Some((client_order_ids, venue_order_ids, pending_op)) = pending else {
940 log::warn!(
941 "Batch order response error without correlation: op={}, ret_code={}, ret_msg={}, req_id={:?}",
942 resp.op,
943 resp.ret_code,
944 resp.ret_msg,
945 resp.req_id,
946 );
947 return;
948 };
949
950 for (index, client_order_id) in client_order_ids.into_iter().enumerate() {
951 let Some(identity) = state
952 .order_identities
953 .get(&client_order_id)
954 .map(|identity| identity.clone())
955 else {
956 log::warn!(
957 "Batch order response error for untracked order: op={}, client_order_id={client_order_id}, ret_msg={}",
958 resp.op,
959 resp.ret_msg,
960 );
961 continue;
962 };
963 emit_rejection_for_op(
964 &pending_op,
965 client_order_id,
966 &identity,
967 venue_order_ids.get(index).copied().flatten(),
968 &resp.ret_msg,
969 emitter,
970 state,
971 ts_init,
972 );
973 }
974 return;
975 }
976
977 if is_bybit_ambiguous_order_error_code(resp.ret_code) {
978 log::warn!(
979 "Ambiguous order response failure: op={}, ret_code={}, ret_msg={}; awaiting reconciliation",
980 resp.op,
981 resp.ret_code,
982 resp.ret_msg,
983 );
984 return;
985 }
986
987 let effective_op = pending
988 .as_ref()
989 .map(|(_, _, op)| *op)
990 .or_else(|| pending_op_from_str(resp.op.as_str()))
991 .unwrap_or_else(|| {
992 log::warn!("Unknown order operation '{}', defaulting to Place", resp.op);
993 PendingOperation::Place
994 });
995
996 let client_order_id = extract_order_link_id_from_data(&resp.data).or_else(|| {
998 pending
999 .as_ref()
1000 .and_then(|(cids, _, _)| cids.first().copied())
1001 });
1002
1003 let stored_venue_order_id = pending
1004 .as_ref()
1005 .and_then(|(_, voids, _)| voids.first().and_then(|v| *v));
1006
1007 let Some(client_order_id) = client_order_id else {
1008 log::warn!(
1009 "Order response error without correlation: op={}, ret_code={}, ret_msg={}, req_id={:?}",
1010 resp.op,
1011 resp.ret_code,
1012 resp.ret_msg,
1013 resp.req_id,
1014 );
1015 return;
1016 };
1017 let Some(identity) = state
1018 .order_identities
1019 .get(&client_order_id)
1020 .map(|r| r.clone())
1021 else {
1022 log::warn!(
1023 "Order response error for untracked order: op={}, client_order_id={client_order_id}, ret_msg={}",
1024 resp.op,
1025 resp.ret_msg,
1026 );
1027 return;
1028 };
1029
1030 let venue_order_id = extract_venue_order_id_from_data(&resp.data).or(stored_venue_order_id);
1031
1032 emit_rejection_for_op(
1033 &effective_op,
1034 client_order_id,
1035 &identity,
1036 venue_order_id,
1037 &resp.ret_msg,
1038 emitter,
1039 state,
1040 ts_init,
1041 );
1042}
1043
1044#[expect(clippy::too_many_arguments)]
1046fn emit_rejection_for_op(
1047 pending_op: &PendingOperation,
1048 client_order_id: ClientOrderId,
1049 identity: &OrderIdentity,
1050 venue_order_id: Option<VenueOrderId>,
1051 reason: &str,
1052 emitter: &ExecutionEventEmitter,
1053 state: &WsDispatchState,
1054 ts_init: UnixNanos,
1055) {
1056 match pending_op {
1057 PendingOperation::Place => {
1058 state.order_identities.remove(&client_order_id);
1059 state.order_snapshots.remove(&client_order_id);
1060 emitter.emit_order_rejected_event(
1061 identity.strategy_id,
1062 identity.instrument_id,
1063 client_order_id,
1064 reason,
1065 ts_init,
1066 false,
1067 );
1068 }
1069 PendingOperation::Cancel => {
1070 emitter.emit_order_cancel_rejected_event(
1071 identity.strategy_id,
1072 identity.instrument_id,
1073 client_order_id,
1074 venue_order_id,
1075 reason,
1076 ts_init,
1077 );
1078 }
1079 PendingOperation::Amend => {
1080 emitter.emit_order_modify_rejected_event(
1081 identity.strategy_id,
1082 identity.instrument_id,
1083 client_order_id,
1084 venue_order_id,
1085 reason,
1086 ts_init,
1087 );
1088 }
1089 }
1090}
1091
1092fn pending_op_from_str(op: &str) -> Option<PendingOperation> {
1094 match op {
1095 BYBIT_OP_ORDER_CREATE => Some(PendingOperation::Place),
1096 BYBIT_OP_ORDER_CANCEL => Some(PendingOperation::Cancel),
1097 BYBIT_OP_ORDER_AMEND => Some(PendingOperation::Amend),
1098 _ => None,
1099 }
1100}
1101
1102fn is_batch_order_op(op: &str) -> bool {
1103 matches!(
1104 op,
1105 BYBIT_OP_ORDER_CREATE_BATCH | BYBIT_OP_ORDER_AMEND_BATCH | BYBIT_OP_ORDER_CANCEL_BATCH
1106 )
1107}
1108
1109fn parse_order_snapshot(
1111 order: &BybitWsAccountOrder,
1112 instrument: &InstrumentAny,
1113) -> Option<OrderStateSnapshot> {
1114 let quantity =
1115 parse_quantity_with_precision(&order.qty, instrument.size_precision(), "order.qty").ok()?;
1116
1117 let price = if !order.price.is_empty() && order.price != "0" {
1118 parse_price_with_precision(&order.price, instrument.price_precision(), "order.price").ok()
1119 } else {
1120 None
1121 };
1122
1123 let trigger_price = if !order.trigger_price.is_empty() && order.trigger_price != "0" {
1124 parse_price_with_precision(
1125 &order.trigger_price,
1126 instrument.price_precision(),
1127 "order.triggerPrice",
1128 )
1129 .ok()
1130 } else {
1131 None
1132 };
1133
1134 Some(OrderStateSnapshot {
1135 quantity,
1136 price,
1137 trigger_price,
1138 })
1139}
1140
1141fn is_snapshot_updated(
1143 snapshot: &OrderStateSnapshot,
1144 client_order_id: &ClientOrderId,
1145 state: &WsDispatchState,
1146) -> bool {
1147 let Some(previous) = state.order_snapshots.get(client_order_id) else {
1148 return false;
1149 };
1150
1151 if let (Some(prev_price), Some(new_price)) = (previous.price, snapshot.price)
1152 && prev_price != new_price
1153 {
1154 return true;
1155 }
1156
1157 if let (Some(prev_trigger), Some(new_trigger)) =
1158 (previous.trigger_price, snapshot.trigger_price)
1159 && prev_trigger != new_trigger
1160 {
1161 return true;
1162 }
1163
1164 previous.quantity != snapshot.quantity
1165}
1166
1167fn ensure_accepted_emitted(
1170 client_order_id: ClientOrderId,
1171 account_id: AccountId,
1172 venue_order_id: VenueOrderId,
1173 identity: &OrderIdentity,
1174 emitter: &ExecutionEventEmitter,
1175 state: &WsDispatchState,
1176 ts_init: UnixNanos,
1177) {
1178 if state.emitted_accepted.contains(&client_order_id) {
1179 return;
1180 }
1181 state.insert_accepted(client_order_id);
1182 let accepted = OrderAccepted::new(
1183 emitter.trader_id(),
1184 identity.strategy_id,
1185 identity.instrument_id,
1186 client_order_id,
1187 venue_order_id,
1188 account_id,
1189 UUID4::new(),
1190 ts_init,
1191 ts_init,
1192 false,
1193 );
1194 emitter.send_order_event(OrderEventAny::Accepted(accepted));
1195}
1196
1197fn cleanup_terminal(client_order_id: ClientOrderId, state: &WsDispatchState) {
1199 state.order_identities.remove(&client_order_id);
1200 state.order_snapshots.remove(&client_order_id);
1201 state.emitted_accepted.remove(&client_order_id);
1202 state.triggered_orders.remove(&client_order_id);
1203 state.filled_orders.remove(&client_order_id);
1204 state.spot_repay_fills.remove(&client_order_id);
1205}
1206
1207fn extract_order_link_id_from_data(data: &serde_json::Value) -> Option<ClientOrderId> {
1209 data.get("orderLinkId")
1210 .and_then(|v| v.as_str())
1211 .filter(|s| !s.is_empty())
1212 .map(ClientOrderId::new)
1213}
1214
1215fn extract_venue_order_id_from_data(data: &serde_json::Value) -> Option<VenueOrderId> {
1217 data.get("orderId")
1218 .and_then(|v| v.as_str())
1219 .filter(|s| !s.is_empty())
1220 .map(VenueOrderId::new)
1221}
1222
1223fn record_spot_repay_fill(
1224 client_order_id: ClientOrderId,
1225 filled: &OrderFilled,
1226 instrument: &InstrumentAny,
1227 state: &WsDispatchState,
1228) {
1229 let Some(base_currency) = instrument.base_currency() else {
1230 return;
1231 };
1232
1233 let base_fee = filled
1234 .commission
1235 .filter(|fee| fee.currency.code == base_currency.code)
1236 .map_or(Decimal::ZERO, |fee| fee.as_decimal().max(Decimal::ZERO));
1237 let fill = SpotRepayFill {
1238 quantity: filled.last_qty,
1239 base_fee,
1240 };
1241 state
1242 .spot_repay_fills
1243 .entry(client_order_id)
1244 .and_modify(|total| {
1245 total.quantity = total.quantity + fill.quantity;
1246 total.base_fee += fill.base_fee;
1247 })
1248 .or_insert(fill);
1249}
1250
1251fn enqueue_spot_repay(
1253 client_order_id: ClientOrderId,
1254 instrument: &InstrumentAny,
1255 state: &WsDispatchState,
1256) {
1257 let Some(base_currency) = instrument.base_currency() else {
1258 return;
1259 };
1260 let Some((_, fill)) = state.spot_repay_fills.remove(&client_order_id) else {
1261 return;
1262 };
1263
1264 state.enqueue_repay(RepayRequest {
1265 coin: base_currency.code,
1266 quantity: fill.quantity,
1267 base_fee: fill.base_fee,
1268 repayment_precision: fill.quantity.precision.max(base_currency.precision),
1269 });
1270}
1271
1272#[cfg(test)]
1273mod tests {
1274 use ahash::AHashMap;
1275 use nautilus_common::messages::{ExecutionEvent, execution::ExecutionReport};
1276 use nautilus_core::{
1277 UnixNanos,
1278 time::{AtomicTime, get_atomic_clock_realtime},
1279 };
1280 use nautilus_live::emitter::ExecutionEventEmitter;
1281 use nautilus_model::{
1282 enums::{AccountType, OrderSide, OrderType},
1283 events::OrderEventAny,
1284 identifiers::{
1285 AccountId, ClientOrderId, InstrumentId, PositionId, StrategyId, TraderId, VenueOrderId,
1286 },
1287 instruments::{Instrument, InstrumentAny},
1288 };
1289 use rstest::rstest;
1290 use ustr::Ustr;
1291
1292 use super::*;
1293 use crate::{
1294 common::{
1295 enums::{BybitExecType, BybitOrderSide, BybitProductType},
1296 parse::{parse_linear_instrument, parse_spot_instrument},
1297 testing::load_test_json,
1298 },
1299 http::models::{BybitFeeRate, BybitInstrumentLinearResponse, BybitInstrumentSpotResponse},
1300 websocket::messages::{
1301 BybitWsAccountExecutionFastMsg, BybitWsMessage, BybitWsOrderResponse,
1302 },
1303 };
1304
1305 fn sample_fee_rate(
1306 symbol: &str,
1307 taker: &str,
1308 maker: &str,
1309 base_coin: Option<&str>,
1310 ) -> BybitFeeRate {
1311 BybitFeeRate {
1312 symbol: Ustr::from(symbol),
1313 taker_fee_rate: taker.to_string(),
1314 maker_fee_rate: maker.to_string(),
1315 base_coin: base_coin.map(Ustr::from),
1316 }
1317 }
1318
1319 fn linear_instrument() -> InstrumentAny {
1320 let json = load_test_json("http_get_instruments_linear.json");
1321 let response: BybitInstrumentLinearResponse = serde_json::from_str(&json).unwrap();
1322 let instrument = &response.result.list[0];
1323 let fee_rate = sample_fee_rate("BTCUSDT", "0.00055", "0.0001", Some("BTC"));
1324 let ts = UnixNanos::new(1_700_000_000_000_000_000);
1325 parse_linear_instrument(instrument, &fee_rate, ts, ts).unwrap()
1326 }
1327
1328 fn spot_instrument() -> InstrumentAny {
1329 let json = load_test_json("http_get_instruments_spot.json");
1330 let response: BybitInstrumentSpotResponse = serde_json::from_str(&json).unwrap();
1331 let instrument = &response.result.list[0];
1332 let fee_rate = sample_fee_rate("BTCUSDT", "0.0006", "0.0001", Some("BTC"));
1333 let ts = UnixNanos::new(1_700_000_000_000_000_000);
1334 parse_spot_instrument(instrument, &fee_rate, ts, ts).unwrap()
1335 }
1336
1337 fn build_instruments(instruments: &[InstrumentAny]) -> AHashMap<Ustr, InstrumentAny> {
1338 let mut map = AHashMap::new();
1339 for inst in instruments {
1340 map.insert(inst.id().symbol.inner(), inst.clone());
1341 }
1342 map
1343 }
1344
1345 fn test_account_id() -> AccountId {
1346 AccountId::from("BYBIT-001")
1347 }
1348
1349 fn create_emitter() -> (
1350 ExecutionEventEmitter,
1351 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
1352 ) {
1353 let clock = get_atomic_clock_realtime();
1354 let trader_id = TraderId::from("TESTER-001");
1355 let account_id = test_account_id();
1356 let mut emitter =
1357 ExecutionEventEmitter::new(clock, trader_id, account_id, AccountType::Margin, None);
1358 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1359 emitter.set_sender(tx);
1360 (emitter, rx)
1361 }
1362
1363 fn default_identity() -> OrderIdentity {
1364 OrderIdentity {
1365 instrument_id: InstrumentId::from("BTCUSDT-LINEAR.BYBIT"),
1366 strategy_id: StrategyId::from("S-001"),
1367 order_side: OrderSide::Buy,
1368 order_type: OrderType::Limit,
1369 venue_position_id: None,
1370 }
1371 }
1372
1373 #[rstest]
1374 #[case::base_fee("BTC", "0.0000015", "0.0000025", "0.000004")]
1375 #[case::base_rebate("BTC", "-0.0000015", "-0.0000025", "0")]
1376 #[case::quote_fee("USDT", "0.075", "0.125", "0")]
1377 fn test_spot_repay_uses_accumulated_execution_quantity(
1378 #[case] fee_currency: &str,
1379 #[case] first_fee: &str,
1380 #[case] second_fee: &str,
1381 #[case] expected_base_fee: &str,
1382 ) {
1383 let instrument = spot_instrument();
1384 let instruments = build_instruments(std::slice::from_ref(&instrument));
1385 let (emitter, _rx) = create_emitter();
1386 let clock = get_atomic_clock_realtime();
1387 let state = WsDispatchState::default();
1388 let (repay_tx, mut repay_rx) = tokio::sync::mpsc::unbounded_channel();
1389 state.set_repay_sender(repay_tx);
1390
1391 let json = load_test_json("ws_account_execution.json");
1392 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
1393 value["data"][0]["category"] = serde_json::Value::String("spot".to_string());
1394 value["data"][0]["symbol"] = serde_json::Value::String("BTCUSDT".to_string());
1395 value["data"][0]["side"] = serde_json::Value::String("Buy".to_string());
1396 value["data"][0]["orderType"] = serde_json::Value::String("Market".to_string());
1397 value["data"][0]["orderQty"] = serde_json::Value::String("100".to_string());
1398 value["data"][0]["execQty"] = serde_json::Value::String("0.0015".to_string());
1399 value["data"][0]["leavesQty"] = serde_json::Value::String("0.0025".to_string());
1400 value["data"][0]["execFee"] = serde_json::Value::String(first_fee.to_string());
1401 value["data"][0]["feeCurrency"] = serde_json::Value::String(fee_currency.to_string());
1402
1403 let first: crate::websocket::messages::BybitWsAccountExecutionMsg =
1404 serde_json::from_value(value.clone()).unwrap();
1405 let client_order_id = ClientOrderId::new(first.data[0].order_link_id);
1406 state.order_identities.insert(
1407 client_order_id,
1408 OrderIdentity {
1409 instrument_id: InstrumentId::from("BTCUSDT-SPOT.BYBIT"),
1410 order_type: OrderType::Market,
1411 ..default_identity()
1412 },
1413 );
1414
1415 dispatch_ws_message(
1416 &BybitWsMessage::AccountExecution(first),
1417 &emitter,
1418 &state,
1419 test_account_id(),
1420 &instruments,
1421 clock,
1422 );
1423 assert!(repay_rx.try_recv().is_err());
1424
1425 value["data"][0]["execId"] = serde_json::Value::String("second-execution".to_string());
1426 value["data"][0]["execQty"] = serde_json::Value::String("0.0025".to_string());
1427 value["data"][0]["leavesQty"] = serde_json::Value::String("0".to_string());
1428 value["data"][0]["execFee"] = serde_json::Value::String(second_fee.to_string());
1429 let second: crate::websocket::messages::BybitWsAccountExecutionMsg =
1430 serde_json::from_value(value).unwrap();
1431
1432 dispatch_ws_message(
1433 &BybitWsMessage::AccountExecution(second),
1434 &emitter,
1435 &state,
1436 test_account_id(),
1437 &instruments,
1438 clock,
1439 );
1440
1441 let repay = repay_rx.try_recv().expect("expected a repay request");
1442 assert_eq!(repay.coin, "BTC");
1443 assert_eq!(repay.quantity, Quantity::from("0.0040"));
1444 assert_eq!(
1445 repay.base_fee,
1446 expected_base_fee.parse::<Decimal>().unwrap()
1447 );
1448 assert_eq!(repay.repayment_precision, 8);
1449 assert!(repay_rx.try_recv().is_err());
1450 }
1451
1452 #[rstest]
1453 fn test_dispatch_tracked_canceled_order_emits_accepted_then_canceled() {
1454 let instrument = linear_instrument();
1455 let instruments = build_instruments(std::slice::from_ref(&instrument));
1456 let (emitter, mut rx) = create_emitter();
1457 let clock = get_atomic_clock_realtime();
1458 let state = WsDispatchState::default();
1459
1460 let json = load_test_json("ws_account_order.json");
1462 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
1463 serde_json::from_str(&json).unwrap();
1464
1465 if let Some(order) = msg.data.first()
1466 && !order.order_link_id.is_empty()
1467 {
1468 let cid = ClientOrderId::new(order.order_link_id);
1469 state.order_identities.insert(cid, default_identity());
1470 }
1471
1472 let ws_msg = BybitWsMessage::AccountOrder(msg);
1473 dispatch_ws_message(
1474 &ws_msg,
1475 &emitter,
1476 &state,
1477 test_account_id(),
1478 &instruments,
1479 clock,
1480 );
1481
1482 let event1 = rx.try_recv().unwrap();
1484 assert!(
1485 matches!(event1, ExecutionEvent::Order(OrderEventAny::Accepted(ref a)) if a.strategy_id == StrategyId::from("S-001")),
1486 "Expected Accepted, found {event1:?}"
1487 );
1488
1489 let event2 = rx.try_recv().unwrap();
1491 assert!(
1492 matches!(event2, ExecutionEvent::Order(OrderEventAny::Canceled(_))),
1493 "Expected Canceled, found {event2:?}"
1494 );
1495 }
1496
1497 #[rstest]
1498 fn test_dispatch_tracked_post_only_cancel_emits_rejected() {
1499 const BYBIT_POST_ONLY_REJECT_REASON: &str = "EC_PostOnlyWillTakeLiquidity";
1500
1501 let instrument = linear_instrument();
1502 let instruments = build_instruments(std::slice::from_ref(&instrument));
1503 let (emitter, mut rx) = create_emitter();
1504 let clock = get_atomic_clock_realtime();
1505 let state = WsDispatchState::default();
1506
1507 let json = load_test_json("ws_account_order.json");
1510 let mut msg: crate::websocket::messages::BybitWsAccountOrderMsg =
1511 serde_json::from_str(&json).unwrap();
1512
1513 let order = msg.data.first_mut().expect("fixture has an order");
1514 order.reject_reason = Ustr::from(BYBIT_POST_ONLY_REJECT_REASON);
1515 order.cum_exec_qty = "0".to_string();
1516 let cid = ClientOrderId::new(order.order_link_id);
1517 state.order_identities.insert(cid, default_identity());
1518
1519 let ws_msg = BybitWsMessage::AccountOrder(msg);
1520 dispatch_ws_message(
1521 &ws_msg,
1522 &emitter,
1523 &state,
1524 test_account_id(),
1525 &instruments,
1526 clock,
1527 );
1528
1529 let event = rx.try_recv().unwrap();
1530 let ExecutionEvent::Order(OrderEventAny::Rejected(rejected)) = event else {
1531 panic!("Expected Rejected, found {event:?}");
1532 };
1533 assert!(rejected.due_post_only);
1534 assert_eq!(rejected.reason, BYBIT_POST_ONLY_REJECT_REASON);
1535 assert_eq!(rejected.client_order_id, cid);
1536 assert!(rx.try_recv().is_err(), "expected only a single event");
1537 }
1538
1539 #[rstest]
1540 fn test_dispatch_untracked_order_emits_report() {
1541 let instrument = linear_instrument();
1542 let instruments = build_instruments(std::slice::from_ref(&instrument));
1543 let (emitter, mut rx) = create_emitter();
1544 let clock = get_atomic_clock_realtime();
1545 let state = WsDispatchState::default();
1546
1547 let json = load_test_json("ws_account_order.json");
1548 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
1549 serde_json::from_str(&json).unwrap();
1550
1551 let ws_msg = BybitWsMessage::AccountOrder(msg);
1553 dispatch_ws_message(
1554 &ws_msg,
1555 &emitter,
1556 &state,
1557 test_account_id(),
1558 &instruments,
1559 clock,
1560 );
1561
1562 let event = rx.try_recv().unwrap();
1563 assert!(matches!(
1564 event,
1565 ExecutionEvent::Report(ExecutionReport::Order(_))
1566 ));
1567 }
1568
1569 #[rstest]
1570 fn test_dispatch_tracked_execution_emits_order_filled() {
1571 let instrument = linear_instrument();
1572 let instruments = build_instruments(std::slice::from_ref(&instrument));
1573 let (emitter, mut rx) = create_emitter();
1574 let clock = get_atomic_clock_realtime();
1575 let state = WsDispatchState::default();
1576
1577 let json = load_test_json("ws_account_execution.json");
1578 let msg: crate::websocket::messages::BybitWsAccountExecutionMsg =
1579 serde_json::from_str(&json).unwrap();
1580
1581 if let Some(exec) = msg.data.first()
1583 && !exec.order_link_id.is_empty()
1584 {
1585 let cid = ClientOrderId::new(exec.order_link_id);
1586 state.order_identities.insert(cid, default_identity());
1587 }
1588
1589 let ws_msg = BybitWsMessage::AccountExecution(msg);
1590 dispatch_ws_message(
1591 &ws_msg,
1592 &emitter,
1593 &state,
1594 test_account_id(),
1595 &instruments,
1596 clock,
1597 );
1598
1599 let event1 = rx.try_recv().unwrap();
1601 assert!(
1602 matches!(event1, ExecutionEvent::Order(OrderEventAny::Accepted(_))),
1603 "Expected Accepted, found {event1:?}"
1604 );
1605
1606 let event2 = rx.try_recv().unwrap();
1608 match event2 {
1609 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
1610 assert_eq!(filled.strategy_id, StrategyId::from("S-001"));
1611 assert_eq!(filled.order_side, OrderSide::Buy);
1612 assert_eq!(filled.order_type, OrderType::Limit);
1613 }
1614 other => panic!("Expected Filled event, found {other:?}"),
1615 }
1616 }
1617
1618 #[rstest]
1619 fn parse_order_filled_uses_payload_fee_currency() {
1620 let instrument = linear_instrument();
1621 let (emitter, _rx) = create_emitter();
1622
1623 let json = load_test_json("ws_account_execution.json");
1624 let msg: crate::websocket::messages::BybitWsAccountExecutionMsg =
1625 serde_json::from_str(&json).unwrap();
1626
1627 let mut exec = msg.data[0].clone();
1628 exec.fee_currency = Ustr::from("BTC");
1629
1630 let filled = parse_order_filled(
1631 &exec,
1632 &instrument,
1633 &default_identity(),
1634 &emitter,
1635 test_account_id(),
1636 UnixNanos::default(),
1637 )
1638 .unwrap();
1639
1640 let commission = filled.commission.expect("commission present");
1641 assert_eq!(commission.currency.code, "BTC");
1642 }
1643
1644 #[rstest]
1645 fn test_dispatch_tracked_execution_preserves_venue_position_id() {
1646 let instrument = linear_instrument();
1647 let instruments = build_instruments(std::slice::from_ref(&instrument));
1648 let (emitter, mut rx) = create_emitter();
1649 let clock = get_atomic_clock_realtime();
1650 let state = WsDispatchState::default();
1651
1652 let json = load_test_json("ws_account_execution.json");
1653 let msg: crate::websocket::messages::BybitWsAccountExecutionMsg =
1654 serde_json::from_str(&json).unwrap();
1655 let venue_position_id = PositionId::from("BTCUSDT-LINEAR.BYBIT-LONG");
1656
1657 if let Some(exec) = msg.data.first()
1658 && !exec.order_link_id.is_empty()
1659 {
1660 let cid = ClientOrderId::new(exec.order_link_id);
1661 state.order_identities.insert(
1662 cid,
1663 OrderIdentity {
1664 venue_position_id: Some(venue_position_id),
1665 ..default_identity()
1666 },
1667 );
1668 }
1669
1670 let ws_msg = BybitWsMessage::AccountExecution(msg);
1671 dispatch_ws_message(
1672 &ws_msg,
1673 &emitter,
1674 &state,
1675 test_account_id(),
1676 &instruments,
1677 clock,
1678 );
1679
1680 let _accepted = rx.try_recv().unwrap();
1681 let event = rx.try_recv().unwrap();
1682 match event {
1683 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
1684 assert_eq!(filled.position_id, Some(venue_position_id));
1685 }
1686 other => panic!("Expected Filled event, found {other:?}"),
1687 }
1688 }
1689
1690 fn fast_execution_msg(is_maker: bool, order_link_id: &str) -> BybitWsAccountExecutionFastMsg {
1691 BybitWsAccountExecutionFastMsg {
1692 topic: Ustr::from("execution.fast"),
1693 id: String::new(),
1694 creation_time: 1_716_800_399_338,
1695 data: vec![BybitWsAccountExecutionFast {
1696 category: BybitProductType::Linear,
1697 symbol: Ustr::from("BTCUSDT"),
1698 exec_id: "fast-1".to_string(),
1699 exec_price: "50000.0".to_string(),
1700 exec_qty: "0.5".to_string(),
1701 order_id: Ustr::from("ord-1"),
1702 order_link_id: Ustr::from(order_link_id),
1703 side: BybitOrderSide::Buy,
1704 exec_time: "1716800399334".to_string(),
1705 is_maker,
1706 seq: 42,
1707 }],
1708 }
1709 }
1710
1711 #[rstest]
1712 fn test_dispatch_tracked_fast_execution_preserves_venue_position_id() {
1713 let instrument = linear_instrument();
1714 let instruments = build_instruments(std::slice::from_ref(&instrument));
1715 let (emitter, mut rx) = create_emitter();
1716 let clock = get_atomic_clock_realtime();
1717 let state = WsDispatchState::default();
1718 let venue_position_id = PositionId::from("BTCUSDT-LINEAR.BYBIT-LONG");
1719
1720 let msg = fast_execution_msg(false, "link-1");
1722 let cid = ClientOrderId::new(msg.data[0].order_link_id);
1723 state.order_identities.insert(
1724 cid,
1725 OrderIdentity {
1726 venue_position_id: Some(venue_position_id),
1727 ..default_identity()
1728 },
1729 );
1730
1731 let ws_msg = BybitWsMessage::AccountExecutionFast(msg);
1732 dispatch_ws_message(
1733 &ws_msg,
1734 &emitter,
1735 &state,
1736 test_account_id(),
1737 &instruments,
1738 clock,
1739 );
1740
1741 let event1 = rx.try_recv().unwrap();
1743 assert!(
1744 matches!(event1, ExecutionEvent::Order(OrderEventAny::Accepted(_))),
1745 "Expected Accepted, found {event1:?}",
1746 );
1747
1748 let event2 = rx.try_recv().unwrap();
1750 match event2 {
1751 ExecutionEvent::Report(ExecutionReport::Fill(report)) => {
1752 assert_eq!(report.venue_position_id, Some(venue_position_id));
1753 assert_eq!(report.client_order_id, Some(cid));
1754 }
1755 other => panic!("Expected FillReport, found {other:?}"),
1756 }
1757 }
1758
1759 #[rstest]
1760 fn test_dispatch_untracked_fast_execution_emits_fill_report_without_position() {
1761 let instrument = linear_instrument();
1762 let instruments = build_instruments(std::slice::from_ref(&instrument));
1763 let (emitter, mut rx) = create_emitter();
1764 let clock = get_atomic_clock_realtime();
1765 let state = WsDispatchState::default();
1766
1767 let msg = fast_execution_msg(true, "");
1769 let ws_msg = BybitWsMessage::AccountExecutionFast(msg);
1770 dispatch_ws_message(
1771 &ws_msg,
1772 &emitter,
1773 &state,
1774 test_account_id(),
1775 &instruments,
1776 clock,
1777 );
1778
1779 let event = rx.try_recv().unwrap();
1781 match event {
1782 ExecutionEvent::Report(ExecutionReport::Fill(report)) => {
1783 assert_eq!(report.client_order_id, None);
1784 assert_eq!(report.venue_position_id, None);
1785 assert_eq!(report.liquidity_side, LiquiditySide::Maker);
1786 assert_eq!(report.commission.as_f64(), 0.0);
1787 }
1788 other => panic!("Expected FillReport, found {other:?}"),
1789 }
1790 }
1791
1792 #[rstest]
1793 fn test_dispatch_untracked_execution_emits_fill_report() {
1794 let instrument = linear_instrument();
1795 let instruments = build_instruments(std::slice::from_ref(&instrument));
1796 let (emitter, mut rx) = create_emitter();
1797 let clock = get_atomic_clock_realtime();
1798 let state = WsDispatchState::default();
1799
1800 let json = load_test_json("ws_account_execution.json");
1801 let msg: crate::websocket::messages::BybitWsAccountExecutionMsg =
1802 serde_json::from_str(&json).unwrap();
1803
1804 let ws_msg = BybitWsMessage::AccountExecution(msg);
1806 dispatch_ws_message(
1807 &ws_msg,
1808 &emitter,
1809 &state,
1810 test_account_id(),
1811 &instruments,
1812 clock,
1813 );
1814
1815 let event = rx.try_recv().unwrap();
1816 assert!(matches!(
1817 event,
1818 ExecutionEvent::Report(ExecutionReport::Fill(_))
1819 ));
1820 }
1821
1822 #[rstest]
1823 #[case::untracked("")]
1824 #[case::tracked("test-order-link-001")]
1825 fn test_dispatch_funding_execution_emits_no_event(#[case] order_link_id: &str) {
1826 let instrument = linear_instrument();
1827 let instruments = build_instruments(std::slice::from_ref(&instrument));
1828 let (emitter, mut rx) = create_emitter();
1829 let clock = get_atomic_clock_realtime();
1830 let state = WsDispatchState::default();
1831
1832 let json = load_test_json("ws_account_execution_funding.json");
1833 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
1834 value["data"][0]["orderLinkId"] = serde_json::Value::String(order_link_id.to_string());
1835 let msg: crate::websocket::messages::BybitWsAccountExecutionMsg =
1836 serde_json::from_value(value).unwrap();
1837 let execution = &msg.data[0];
1838
1839 assert_eq!(execution.exec_type, BybitExecType::Funding);
1840
1841 if !execution.order_link_id.is_empty() {
1842 state.order_identities.insert(
1843 ClientOrderId::new(execution.order_link_id),
1844 default_identity(),
1845 );
1846 }
1847
1848 dispatch_ws_message(
1849 &BybitWsMessage::AccountExecution(msg),
1850 &emitter,
1851 &state,
1852 test_account_id(),
1853 &instruments,
1854 clock,
1855 );
1856
1857 assert!(matches!(
1858 rx.try_recv(),
1859 Err(tokio::sync::mpsc::error::TryRecvError::Empty)
1860 ));
1861 }
1862
1863 #[rstest]
1864 #[case::corporate_action("CorporateAction")]
1865 #[case::forward_split_settle("ForwardSplitSettle")]
1866 #[case::reverse_split_settle("ReverseSplitSettle")]
1867 #[case::dividend("Dividend")]
1868 #[case::unrecognized("StockMerger")]
1869 fn test_dispatch_venue_initiated_execution_emits_only_fill_report(#[case] exec_type: &str) {
1870 let instrument = linear_instrument();
1871 let instruments = build_instruments(std::slice::from_ref(&instrument));
1872 let (emitter, mut rx) = create_emitter();
1873 let clock = get_atomic_clock_realtime();
1874 let state = WsDispatchState::default();
1875
1876 let json = load_test_json("ws_account_execution_adl.json");
1877 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
1878 value["data"][0]["execType"] = serde_json::Value::String(exec_type.to_string());
1879 let msg: crate::websocket::messages::BybitWsAccountExecutionMsg =
1880 serde_json::from_value(value).unwrap();
1881 let execution = &msg.data[0];
1882
1883 assert!(execution.order_link_id.is_empty());
1884
1885 dispatch_ws_message(
1886 &BybitWsMessage::AccountExecution(msg),
1887 &emitter,
1888 &state,
1889 test_account_id(),
1890 &instruments,
1891 clock,
1892 );
1893
1894 let event = rx.try_recv().unwrap();
1895 match event {
1896 ExecutionEvent::Report(ExecutionReport::Fill(report)) => {
1897 assert_eq!(report.client_order_id, None);
1898 assert_eq!(
1899 report.venue_order_id,
1900 VenueOrderId::from("9aac161b-8ed6-450d-9cab-c5cc67c21785")
1901 );
1902 }
1903 other => panic!("Expected FillReport, found {other:?}"),
1904 }
1905 assert!(rx.try_recv().is_err());
1906 }
1907
1908 #[rstest]
1909 fn test_dispatch_wallet_emits_account_state() {
1910 let instruments = AHashMap::new();
1911 let (emitter, mut rx) = create_emitter();
1912 let clock = get_atomic_clock_realtime();
1913 let state = WsDispatchState::default();
1914
1915 let json = load_test_json("ws_account_wallet.json");
1916 let msg: crate::websocket::messages::BybitWsAccountWalletMsg =
1917 serde_json::from_str(&json).unwrap();
1918 let ws_msg = BybitWsMessage::AccountWallet(msg);
1919
1920 dispatch_ws_message(
1921 &ws_msg,
1922 &emitter,
1923 &state,
1924 test_account_id(),
1925 &instruments,
1926 clock,
1927 );
1928
1929 let event = rx.try_recv().unwrap();
1930 assert!(matches!(event, ExecutionEvent::Account(_)));
1931 }
1932
1933 #[rstest]
1934 fn test_dispatch_data_message_ignored() {
1935 let instruments = AHashMap::new();
1936 let (emitter, mut rx) = create_emitter();
1937 let clock = get_atomic_clock_realtime();
1938 let state = WsDispatchState::default();
1939
1940 let json = load_test_json("ws_public_trade.json");
1941 let msg: crate::websocket::messages::BybitWsTradeMsg = serde_json::from_str(&json).unwrap();
1942 let ws_msg = BybitWsMessage::Trade(msg);
1943
1944 dispatch_ws_message(
1945 &ws_msg,
1946 &emitter,
1947 &state,
1948 test_account_id(),
1949 &instruments,
1950 clock,
1951 );
1952
1953 rx.try_recv().unwrap_err();
1954 }
1955
1956 #[rstest]
1957 fn test_accepted_dedup_prevents_duplicate() {
1958 let instrument = linear_instrument();
1959 let instruments = build_instruments(std::slice::from_ref(&instrument));
1960 let (emitter, mut rx) = create_emitter();
1961 let clock = get_atomic_clock_realtime();
1962 let state = WsDispatchState::default();
1963
1964 let json = load_test_json("ws_account_order.json");
1966 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
1967 value["data"][0]["orderStatus"] = serde_json::Value::String("New".to_string());
1968 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
1969 serde_json::from_value(value).unwrap();
1970
1971 if let Some(order) = msg.data.first()
1972 && !order.order_link_id.is_empty()
1973 {
1974 let cid = ClientOrderId::new(order.order_link_id);
1975 state.order_identities.insert(cid, default_identity());
1976 }
1977
1978 let ws_msg = BybitWsMessage::AccountOrder(msg.clone());
1979 dispatch_ws_message(
1980 &ws_msg,
1981 &emitter,
1982 &state,
1983 test_account_id(),
1984 &instruments,
1985 clock,
1986 );
1987
1988 let event = rx.try_recv().unwrap();
1989 assert!(matches!(
1990 event,
1991 ExecutionEvent::Order(OrderEventAny::Accepted(_))
1992 ));
1993
1994 let ws_msg2 = BybitWsMessage::AccountOrder(msg);
1996 dispatch_ws_message(
1997 &ws_msg2,
1998 &emitter,
1999 &state,
2000 test_account_id(),
2001 &instruments,
2002 clock,
2003 );
2004
2005 rx.try_recv().unwrap_err();
2006 }
2007
2008 fn new_order_value() -> serde_json::Value {
2009 let json = load_test_json("ws_account_order.json");
2010 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
2011 value["data"][0]["orderStatus"] = serde_json::Value::String("New".to_string());
2012 value
2013 }
2014
2015 struct DispatchTestContext {
2016 instruments: AHashMap<Ustr, InstrumentAny>,
2017 emitter: ExecutionEventEmitter,
2018 rx: tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
2019 clock: &'static AtomicTime,
2020 state: WsDispatchState,
2021 }
2022
2023 impl DispatchTestContext {
2024 fn new() -> Self {
2025 let instrument = linear_instrument();
2026 let instruments = build_instruments(std::slice::from_ref(&instrument));
2027 let (emitter, rx) = create_emitter();
2028 let clock = get_atomic_clock_realtime();
2029 let state = WsDispatchState::default();
2030 Self {
2031 instruments,
2032 emitter,
2033 rx,
2034 clock,
2035 state,
2036 }
2037 }
2038
2039 fn accept_order(&mut self, value: &serde_json::Value) {
2040 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
2041 serde_json::from_value(value.clone()).unwrap();
2042
2043 if let Some(order) = msg.data.first()
2044 && !order.order_link_id.is_empty()
2045 && !self
2046 .state
2047 .order_identities
2048 .contains_key(&ClientOrderId::new(order.order_link_id))
2049 {
2050 let cid = ClientOrderId::new(order.order_link_id);
2051 self.state.order_identities.insert(cid, default_identity());
2052 }
2053
2054 self.dispatch_value(value);
2055
2056 let event = self.rx.try_recv().unwrap();
2057 assert!(
2058 matches!(event, ExecutionEvent::Order(OrderEventAny::Accepted(_))),
2059 "Expected Accepted, found {event:?}"
2060 );
2061 }
2062
2063 fn dispatch_value(&self, value: &serde_json::Value) {
2064 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
2065 serde_json::from_value(value.clone()).unwrap();
2066 let ws_msg = BybitWsMessage::AccountOrder(msg);
2067 dispatch_ws_message(
2068 &ws_msg,
2069 &self.emitter,
2070 &self.state,
2071 test_account_id(),
2072 &self.instruments,
2073 self.clock,
2074 );
2075 }
2076
2077 fn recv_updated(&mut self) -> OrderUpdated {
2078 let event = self.rx.try_recv().unwrap();
2079 match event {
2080 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => updated,
2081 other => panic!("Expected Updated event, found {other:?}"),
2082 }
2083 }
2084 }
2085
2086 #[rstest]
2087 fn test_dispatch_order_updated_on_price_change() {
2088 let mut ctx = DispatchTestContext::new();
2089 let value = new_order_value();
2090 ctx.accept_order(&value);
2091
2092 let mut amended = value;
2093 amended["data"][0]["price"] = serde_json::Value::String("31000".to_string());
2094 ctx.dispatch_value(&amended);
2095
2096 let updated = ctx.recv_updated();
2097 assert_eq!(updated.client_order_id, ClientOrderId::from("client-1"));
2098 assert_eq!(updated.price, Some(Price::from("31000.00")));
2099 assert_eq!(updated.quantity, Quantity::from("0.010"));
2100 assert_eq!(updated.trigger_price, None);
2101 assert!(updated.venue_order_id.is_some());
2102 }
2103
2104 #[rstest]
2105 fn test_dispatch_order_updated_on_quantity_change() {
2106 let mut ctx = DispatchTestContext::new();
2107 let value = new_order_value();
2108 ctx.accept_order(&value);
2109
2110 let mut amended = value;
2111 amended["data"][0]["qty"] = serde_json::Value::String("0.020".to_string());
2112 ctx.dispatch_value(&amended);
2113
2114 let updated = ctx.recv_updated();
2115 assert_eq!(updated.quantity, Quantity::from("0.020"));
2116 assert_eq!(updated.price, Some(Price::from("30000.00")));
2117 }
2118
2119 #[rstest]
2120 fn test_dispatch_order_updated_on_trigger_price_change() {
2121 let mut ctx = DispatchTestContext::new();
2122 let mut value = new_order_value();
2123 value["data"][0]["triggerPrice"] = serde_json::Value::String("29000".to_string());
2124 ctx.accept_order(&value);
2125
2126 let mut amended = value;
2127 amended["data"][0]["triggerPrice"] = serde_json::Value::String("28000".to_string());
2128 ctx.dispatch_value(&amended);
2129
2130 let updated = ctx.recv_updated();
2131 assert_eq!(updated.trigger_price, Some(Price::from("28000.00")));
2132 assert_eq!(updated.price, Some(Price::from("30000.00")));
2133 }
2134
2135 #[rstest]
2136 fn test_dispatch_dedup_suppresses_identical_after_snapshot() {
2137 let mut ctx = DispatchTestContext::new();
2138 let value = new_order_value();
2139 ctx.accept_order(&value);
2140
2141 ctx.dispatch_value(&value);
2142
2143 assert!(
2144 ctx.rx.try_recv().is_err(),
2145 "Expected no event for identical redelivery"
2146 );
2147 }
2148
2149 #[rstest]
2150 fn test_dispatch_order_updated_stores_snapshot_for_subsequent_change() {
2151 let mut ctx = DispatchTestContext::new();
2152 let value = new_order_value();
2153 ctx.accept_order(&value);
2154
2155 let mut amended1 = value.clone();
2156 amended1["data"][0]["price"] = serde_json::Value::String("31000".to_string());
2157 ctx.dispatch_value(&amended1);
2158 let _result = ctx.recv_updated();
2159
2160 let mut amended2 = value;
2161 amended2["data"][0]["price"] = serde_json::Value::String("32000".to_string());
2162 ctx.dispatch_value(&amended2);
2163
2164 let updated = ctx.recv_updated();
2165 assert_eq!(updated.price, Some(Price::from("32000.00")));
2166 }
2167
2168 #[rstest]
2169 fn test_dispatch_accepted_with_seeded_snapshot_emits_updated_for_bbo() {
2170 let mut ctx = DispatchTestContext::new();
2171 let value = new_order_value();
2172 let cid = ClientOrderId::from("client-1");
2173 ctx.state.order_identities.insert(cid, default_identity());
2174 ctx.state.order_snapshots.insert(
2175 cid,
2176 OrderStateSnapshot {
2177 quantity: Quantity::from("0.010"),
2178 price: Some(Price::from("29000.00")),
2179 trigger_price: None,
2180 },
2181 );
2182
2183 ctx.dispatch_value(&value);
2184
2185 let accepted = ctx.rx.try_recv().unwrap();
2186 assert!(
2187 matches!(accepted, ExecutionEvent::Order(OrderEventAny::Accepted(_))),
2188 "Expected Accepted first, found {accepted:?}"
2189 );
2190
2191 let updated = ctx.recv_updated();
2192 assert_eq!(updated.client_order_id, cid);
2193 assert_eq!(updated.price, Some(Price::from("30000.00")));
2194 assert_eq!(updated.quantity, Quantity::from("0.010"));
2195 assert!(updated.venue_order_id.is_some());
2196 }
2197
2198 #[rstest]
2199 fn test_emit_rejection_for_place_clears_snapshot() {
2200 let ctx = DispatchTestContext::new();
2201 let cid = ClientOrderId::from("client-1");
2202 let identity = default_identity();
2203 ctx.state.order_identities.insert(cid, identity.clone());
2204 ctx.state.order_snapshots.insert(
2205 cid,
2206 OrderStateSnapshot {
2207 quantity: Quantity::from("0.010"),
2208 price: Some(Price::from("30000.00")),
2209 trigger_price: None,
2210 },
2211 );
2212
2213 emit_rejection_for_op(
2214 &PendingOperation::Place,
2215 cid,
2216 &identity,
2217 None,
2218 "rejected",
2219 &ctx.emitter,
2220 &ctx.state,
2221 UnixNanos::from(1u64),
2222 );
2223
2224 assert!(!ctx.state.order_identities.contains_key(&cid));
2225 assert!(!ctx.state.order_snapshots.contains_key(&cid));
2226 }
2227
2228 fn order_response(
2229 op: &str,
2230 ret_code: i64,
2231 ret_msg: &str,
2232 req_id: &str,
2233 data: serde_json::Value,
2234 ret_ext_info: Option<serde_json::Value>,
2235 ) -> BybitWsOrderResponse {
2236 BybitWsOrderResponse {
2237 op: Ustr::from(op),
2238 conn_id: Some("test-conn-id".to_string()),
2239 ret_code,
2240 ret_msg: ret_msg.to_string(),
2241 data,
2242 req_id: Some(req_id.to_string()),
2243 header: None,
2244 ret_ext_info,
2245 }
2246 }
2247
2248 #[rstest]
2249 fn test_dispatch_local_not_sent_emits_order_rejected() {
2250 let mut ctx = DispatchTestContext::new();
2251 let cid = ClientOrderId::from("not-sent-1");
2252 ctx.state.order_identities.insert(cid, default_identity());
2253 ctx.state.pending_requests.insert(
2254 "req-not-sent".to_string(),
2255 (vec![cid], vec![None], PendingOperation::Place),
2256 );
2257 let response = order_response(
2258 BYBIT_OP_ORDER_CREATE,
2259 -1,
2260 "Order command was not written",
2261 "req-not-sent",
2262 serde_json::json!({}),
2263 None,
2264 );
2265
2266 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2267
2268 let event = ctx.rx.try_recv().expect("expected OrderRejected event");
2269 assert!(
2270 matches!(event, ExecutionEvent::Order(OrderEventAny::Rejected(ref rejected))
2271 if rejected.client_order_id == cid
2272 && rejected.reason == "Order command was not written"),
2273 "Expected OrderRejected for {cid}, found {event:?}"
2274 );
2275 assert!(!ctx.state.order_identities.contains_key(&cid));
2276 assert!(!ctx.state.pending_requests.contains_key("req-not-sent"));
2277 }
2278
2279 #[rstest]
2280 fn test_dispatch_single_cancel_rejection_emits_cancel_rejected() {
2281 let mut ctx = DispatchTestContext::new();
2282 let cid = ClientOrderId::from("cancel-reject-1");
2283 let venue_order_id = VenueOrderId::from("venue-cancel-1");
2284 ctx.state.order_identities.insert(cid, default_identity());
2285 ctx.state.pending_requests.insert(
2286 "req-cancel".to_string(),
2287 (
2288 vec![cid],
2289 vec![Some(venue_order_id)],
2290 PendingOperation::Cancel,
2291 ),
2292 );
2293
2294 let response = order_response(
2295 BYBIT_OP_ORDER_CANCEL,
2296 110001,
2297 "Order does not exist.",
2298 "req-cancel",
2299 serde_json::json!({
2300 "orderId": venue_order_id.to_string(),
2301 "orderLinkId": cid.to_string(),
2302 }),
2303 None,
2304 );
2305
2306 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2307
2308 let event = ctx.rx.try_recv().expect("expected CancelRejected event");
2309 assert!(
2310 matches!(event, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected)) if rejected.client_order_id == cid),
2311 "Expected CancelRejected for {cid}, found {event:?}"
2312 );
2313 assert!(
2314 ctx.rx.try_recv().is_err(),
2315 "Expected no extra events for single cancel rejection"
2316 );
2317 }
2318
2319 #[rstest]
2320 fn test_dispatch_single_modify_rejection_emits_modify_rejected() {
2321 let mut ctx = DispatchTestContext::new();
2322 let cid = ClientOrderId::from("modify-reject-1");
2323 let venue_order_id = VenueOrderId::from("venue-modify-1");
2324 ctx.state.order_identities.insert(cid, default_identity());
2325 ctx.state.pending_requests.insert(
2326 "req-modify".to_string(),
2327 (
2328 vec![cid],
2329 vec![Some(venue_order_id)],
2330 PendingOperation::Amend,
2331 ),
2332 );
2333
2334 let response = order_response(
2335 BYBIT_OP_ORDER_AMEND,
2336 110003,
2337 "Order price exceeds allowable range.",
2338 "req-modify",
2339 serde_json::json!({
2340 "orderId": venue_order_id.to_string(),
2341 "orderLinkId": cid.to_string(),
2342 }),
2343 None,
2344 );
2345
2346 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2347
2348 let event = ctx.rx.try_recv().expect("expected ModifyRejected event");
2349 assert!(
2350 matches!(event, ExecutionEvent::Order(OrderEventAny::ModifyRejected(ref rejected)) if rejected.client_order_id == cid),
2351 "Expected ModifyRejected for {cid}, found {event:?}"
2352 );
2353 assert!(
2354 ctx.rx.try_recv().is_err(),
2355 "Expected no extra events for single modify rejection"
2356 );
2357 }
2358
2359 #[rstest]
2360 fn test_dispatch_single_cancel_rate_limit_emits_cancel_rejected() {
2361 let mut ctx = DispatchTestContext::new();
2362 let cid = ClientOrderId::from("cancel-rate-limit-1");
2363 let venue_order_id = VenueOrderId::from("venue-cancel-rate-limit-1");
2364 ctx.state.order_identities.insert(cid, default_identity());
2365 ctx.state.pending_requests.insert(
2366 "req-cancel-rate-limit".to_string(),
2367 (
2368 vec![cid],
2369 vec![Some(venue_order_id)],
2370 PendingOperation::Cancel,
2371 ),
2372 );
2373
2374 let response = order_response(
2375 BYBIT_OP_ORDER_CANCEL,
2376 10429,
2377 "System level frequency protection.",
2378 "req-cancel-rate-limit",
2379 serde_json::json!({
2380 "orderId": venue_order_id.to_string(),
2381 "orderLinkId": cid.to_string(),
2382 }),
2383 None,
2384 );
2385
2386 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2387
2388 let event = ctx.rx.try_recv().expect("expected CancelRejected event");
2389 assert!(
2390 matches!(event, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected))
2391 if rejected.client_order_id == cid && rejected.venue_order_id == Some(venue_order_id)),
2392 "Expected CancelRejected for {cid}, found {event:?}"
2393 );
2394 assert!(ctx.state.order_identities.contains_key(&cid));
2395 }
2396
2397 #[rstest]
2398 fn test_dispatch_single_modify_server_error_keeps_outcome_unresolved() {
2399 let mut ctx = DispatchTestContext::new();
2400 let cid = ClientOrderId::from("modify-server-error-1");
2401 let venue_order_id = VenueOrderId::from("venue-modify-server-error-1");
2402 ctx.state.order_identities.insert(cid, default_identity());
2403 ctx.state.pending_requests.insert(
2404 "req-modify-server-error".to_string(),
2405 (
2406 vec![cid],
2407 vec![Some(venue_order_id)],
2408 PendingOperation::Amend,
2409 ),
2410 );
2411
2412 let response = order_response(
2413 BYBIT_OP_ORDER_AMEND,
2414 10016,
2415 "Internal server error.",
2416 "req-modify-server-error",
2417 serde_json::json!({
2418 "orderId": venue_order_id.to_string(),
2419 "orderLinkId": cid.to_string(),
2420 }),
2421 None,
2422 );
2423
2424 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2425
2426 assert!(
2427 ctx.rx.try_recv().is_err(),
2428 "Expected no ModifyRejected event for server error response"
2429 );
2430 assert!(ctx.state.order_identities.contains_key(&cid));
2431 }
2432
2433 #[rstest]
2434 fn test_dispatch_batch_cancel_top_level_failure_does_not_emit_cancel_rejected() {
2435 let mut ctx = DispatchTestContext::new();
2436 let cid_1 = ClientOrderId::from("batch-cancel-1");
2437 let cid_2 = ClientOrderId::from("batch-cancel-2");
2438 ctx.state.order_identities.insert(cid_1, default_identity());
2439 ctx.state.order_identities.insert(cid_2, default_identity());
2440 ctx.state.pending_requests.insert(
2441 "req-batch-cancel".to_string(),
2442 (
2443 vec![cid_1, cid_2],
2444 vec![
2445 Some(VenueOrderId::from("venue-batch-1")),
2446 Some(VenueOrderId::from("venue-batch-2")),
2447 ],
2448 PendingOperation::Cancel,
2449 ),
2450 );
2451
2452 let response = order_response(
2453 BYBIT_OP_ORDER_CANCEL_BATCH,
2454 10016,
2455 "Internal server error.",
2456 "req-batch-cancel",
2457 serde_json::json!({}),
2458 None,
2459 );
2460
2461 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2462
2463 assert!(
2464 ctx.rx.try_recv().is_err(),
2465 "Expected no CancelRejected events for ambiguous batch failure"
2466 );
2467 assert!(ctx.state.order_identities.contains_key(&cid_1));
2468 assert!(ctx.state.order_identities.contains_key(&cid_2));
2469 }
2470
2471 #[rstest]
2472 fn test_dispatch_single_item_batch_cancel_top_level_failure_does_not_emit_cancel_rejected() {
2473 let mut ctx = DispatchTestContext::new();
2474 let cid = ClientOrderId::from("batch-cancel-single");
2475 ctx.state.order_identities.insert(cid, default_identity());
2476 ctx.state.pending_requests.insert(
2477 "req-batch-cancel-single".to_string(),
2478 (
2479 vec![cid],
2480 vec![Some(VenueOrderId::from("venue-batch-single"))],
2481 PendingOperation::Cancel,
2482 ),
2483 );
2484
2485 let response = order_response(
2486 BYBIT_OP_ORDER_CANCEL_BATCH,
2487 10016,
2488 "Internal server error.",
2489 "req-batch-cancel-single",
2490 serde_json::json!({}),
2491 None,
2492 );
2493
2494 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2495
2496 assert!(
2497 ctx.rx.try_recv().is_err(),
2498 "Expected no CancelRejected event for ambiguous one-item batch failure"
2499 );
2500 assert!(ctx.state.order_identities.contains_key(&cid));
2501 }
2502
2503 #[rstest]
2504 fn test_dispatch_batch_cancel_top_level_rate_limit_emits_cancel_rejected() {
2505 let mut ctx = DispatchTestContext::new();
2506 let cid_1 = ClientOrderId::from("batch-rate-1");
2507 let cid_2 = ClientOrderId::from("batch-rate-2");
2508 let venue_1 = VenueOrderId::from("venue-rate-1");
2509 let venue_2 = VenueOrderId::from("venue-rate-2");
2510 ctx.state.order_identities.insert(cid_1, default_identity());
2511 ctx.state.order_identities.insert(cid_2, default_identity());
2512 ctx.state.pending_requests.insert(
2513 "req-batch-rate".to_string(),
2514 (
2515 vec![cid_1, cid_2],
2516 vec![Some(venue_1), Some(venue_2)],
2517 PendingOperation::Cancel,
2518 ),
2519 );
2520
2521 let response = order_response(
2522 BYBIT_OP_ORDER_CANCEL_BATCH,
2523 10006,
2524 "Too many visits.",
2525 "req-batch-rate",
2526 serde_json::json!({}),
2527 None,
2528 );
2529
2530 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2531
2532 let first = ctx.rx.try_recv().expect("expected first CancelRejected");
2533 let second = ctx.rx.try_recv().expect("expected second CancelRejected");
2534 assert!(
2535 matches!(first, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected))
2536 if rejected.client_order_id == cid_1 && rejected.venue_order_id == Some(venue_1))
2537 );
2538 assert!(
2539 matches!(second, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected))
2540 if rejected.client_order_id == cid_2 && rejected.venue_order_id == Some(venue_2))
2541 );
2542 }
2543
2544 #[rstest]
2545 fn test_dispatch_batch_cancel_per_item_error_emits_only_failed_item() {
2546 let mut ctx = DispatchTestContext::new();
2547 let cid_1 = ClientOrderId::from("batch-cancel-ok");
2548 let cid_2 = ClientOrderId::from("batch-cancel-reject");
2549 ctx.state.order_identities.insert(cid_1, default_identity());
2550 ctx.state.order_identities.insert(cid_2, default_identity());
2551 ctx.state.pending_requests.insert(
2552 "req-batch-mixed".to_string(),
2553 (
2554 vec![cid_1, cid_2],
2555 vec![
2556 Some(VenueOrderId::from("venue-batch-ok")),
2557 Some(VenueOrderId::from("venue-batch-reject")),
2558 ],
2559 PendingOperation::Cancel,
2560 ),
2561 );
2562
2563 let response = order_response(
2564 BYBIT_OP_ORDER_CANCEL_BATCH,
2565 0,
2566 "OK",
2567 "req-batch-mixed",
2568 serde_json::json!({
2569 "list": [
2570 {"orderId": "venue-batch-ok", "orderLinkId": cid_1.to_string()},
2571 {"orderId": "venue-batch-reject", "orderLinkId": cid_2.to_string()}
2572 ]
2573 }),
2574 Some(serde_json::json!({
2575 "list": [
2576 {"code": 0, "msg": "OK"},
2577 {"code": 170213, "msg": "Order does not exist."}
2578 ]
2579 })),
2580 );
2581
2582 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2583
2584 let event = ctx.rx.try_recv().expect("expected CancelRejected event");
2585 assert!(
2586 matches!(event, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected)) if rejected.client_order_id == cid_2),
2587 "Expected CancelRejected for failed batch item {cid_2}, found {event:?}"
2588 );
2589 assert!(
2590 ctx.rx.try_recv().is_err(),
2591 "Expected no CancelRejected event for successful batch item"
2592 );
2593 }
2594
2595 #[rstest]
2596 fn test_dispatch_batch_cancel_per_item_rate_limit_emits_cancel_rejected() {
2597 let mut ctx = DispatchTestContext::new();
2598 let cid = ClientOrderId::from("batch-cancel-rate-limit");
2599 let venue_order_id = VenueOrderId::from("venue-batch-rate-limit");
2600 ctx.state.order_identities.insert(cid, default_identity());
2601 ctx.state.pending_requests.insert(
2602 "req-batch-rate-limit".to_string(),
2603 (
2604 vec![cid],
2605 vec![Some(venue_order_id)],
2606 PendingOperation::Cancel,
2607 ),
2608 );
2609
2610 let response = order_response(
2611 BYBIT_OP_ORDER_CANCEL_BATCH,
2612 0,
2613 "OK",
2614 "req-batch-rate-limit",
2615 serde_json::json!({
2616 "list": [
2617 {"orderId": venue_order_id.to_string(), "orderLinkId": cid.to_string()}
2618 ]
2619 }),
2620 Some(serde_json::json!({
2621 "list": [
2622 {"code": 10006, "msg": "Too many visits."}
2623 ]
2624 })),
2625 );
2626
2627 dispatch_order_response(&response, &ctx.emitter, &ctx.state, UnixNanos::from(1u64));
2628
2629 let event = ctx.rx.try_recv().expect("expected CancelRejected event");
2630 assert!(
2631 matches!(event, ExecutionEvent::Order(OrderEventAny::CancelRejected(ref rejected))
2632 if rejected.client_order_id == cid && rejected.venue_order_id == Some(venue_order_id)),
2633 "Expected CancelRejected for {cid}, found {event:?}"
2634 );
2635 assert!(ctx.state.order_identities.contains_key(&cid));
2636 }
2637
2638 #[rstest]
2639 fn test_dispatch_filled_with_seeded_snapshot_emits_updated_for_bbo() {
2640 let mut ctx = DispatchTestContext::new();
2641 let mut value = new_order_value();
2642 value["data"][0]["orderStatus"] = serde_json::Value::String("Filled".to_string());
2643 value["data"][0]["cumExecQty"] = serde_json::Value::String("0.010".to_string());
2644 let cid = ClientOrderId::from("client-1");
2645 ctx.state.order_identities.insert(cid, default_identity());
2646 ctx.state.order_snapshots.insert(
2647 cid,
2648 OrderStateSnapshot {
2649 quantity: Quantity::from("0.010"),
2650 price: Some(Price::from("29000.00")),
2651 trigger_price: None,
2652 },
2653 );
2654
2655 ctx.dispatch_value(&value);
2656
2657 let accepted_event = ctx.rx.try_recv().unwrap();
2658 let accepted = match accepted_event {
2659 ExecutionEvent::Order(OrderEventAny::Accepted(accepted)) => accepted,
2660 other => panic!("Expected Accepted first, found {other:?}"),
2661 };
2662 assert_eq!(accepted.client_order_id, cid);
2663 assert!(!accepted.venue_order_id.to_string().is_empty());
2664
2665 let updated = ctx.recv_updated();
2666 assert_eq!(updated.client_order_id, cid);
2667 assert_eq!(updated.price, Some(Price::from("30000.00")));
2668 assert_eq!(updated.quantity, Quantity::from("0.010"));
2669 assert_eq!(updated.trigger_price, None);
2670 assert_eq!(updated.venue_order_id, Some(accepted.venue_order_id));
2671
2672 let stored = ctx.state.order_snapshots.get(&cid).unwrap();
2673 assert_eq!(stored.price, Some(Price::from("30000.00")));
2674 }
2675
2676 #[rstest]
2677 fn test_dispatch_accepted_with_matching_seeded_snapshot_no_updated() {
2678 let mut ctx = DispatchTestContext::new();
2679 let value = new_order_value();
2680 let cid = ClientOrderId::from("client-1");
2681 ctx.state.order_identities.insert(cid, default_identity());
2682 ctx.state.order_snapshots.insert(
2683 cid,
2684 OrderStateSnapshot {
2685 quantity: Quantity::from("0.010"),
2686 price: Some(Price::from("30000.00")),
2687 trigger_price: None,
2688 },
2689 );
2690
2691 ctx.dispatch_value(&value);
2692
2693 let accepted = ctx.rx.try_recv().unwrap();
2694 assert!(matches!(
2695 accepted,
2696 ExecutionEvent::Order(OrderEventAny::Accepted(_))
2697 ));
2698
2699 assert!(
2700 ctx.rx.try_recv().is_err(),
2701 "Expected no OrderUpdated when seed matches venue snapshot"
2702 );
2703 }
2704
2705 #[rstest]
2706 #[case::price_changed(
2707 Some(Price::from("100.00")),
2708 None,
2709 Quantity::from("1.000"),
2710 Some(Price::from("200.00")),
2711 None,
2712 Quantity::from("1.000"),
2713 true
2714 )]
2715 #[case::trigger_changed(
2716 None,
2717 Some(Price::from("100.00")),
2718 Quantity::from("1.000"),
2719 None,
2720 Some(Price::from("90.00")),
2721 Quantity::from("1.000"),
2722 true
2723 )]
2724 #[case::qty_changed(
2725 Some(Price::from("100.00")),
2726 None,
2727 Quantity::from("1.000"),
2728 Some(Price::from("100.00")),
2729 None,
2730 Quantity::from("2.000"),
2731 true
2732 )]
2733 #[case::no_change(
2734 Some(Price::from("100.00")),
2735 None,
2736 Quantity::from("1.000"),
2737 Some(Price::from("100.00")),
2738 None,
2739 Quantity::from("1.000"),
2740 false
2741 )]
2742 fn test_is_snapshot_updated(
2743 #[case] prev_price: Option<Price>,
2744 #[case] prev_trigger: Option<Price>,
2745 #[case] prev_qty: Quantity,
2746 #[case] new_price: Option<Price>,
2747 #[case] new_trigger: Option<Price>,
2748 #[case] new_qty: Quantity,
2749 #[case] expected: bool,
2750 ) {
2751 let state = WsDispatchState::default();
2752 let cid = ClientOrderId::from("test-1");
2753 state.order_snapshots.insert(
2754 cid,
2755 OrderStateSnapshot {
2756 quantity: prev_qty,
2757 price: prev_price,
2758 trigger_price: prev_trigger,
2759 },
2760 );
2761
2762 let new_snapshot = OrderStateSnapshot {
2763 quantity: new_qty,
2764 price: new_price,
2765 trigger_price: new_trigger,
2766 };
2767 assert_eq!(is_snapshot_updated(&new_snapshot, &cid, &state), expected);
2768 }
2769
2770 #[rstest]
2771 fn test_is_snapshot_updated_no_previous() {
2772 let state = WsDispatchState::default();
2773 let cid = ClientOrderId::from("test-1");
2774
2775 let new_snapshot = OrderStateSnapshot {
2776 quantity: Quantity::from("1.000"),
2777 price: Some(Price::from("100.00")),
2778 trigger_price: None,
2779 };
2780 assert!(!is_snapshot_updated(&new_snapshot, &cid, &state));
2781 }
2782
2783 #[rstest]
2784 #[case::limit_order("30000", "0", Some(Price::from("30000.00")), None)]
2785 #[case::conditional("0", "29000", None, Some(Price::from("29000.00")))]
2786 #[case::both(
2787 "30000",
2788 "29000",
2789 Some(Price::from("30000.00")),
2790 Some(Price::from("29000.00"))
2791 )]
2792 fn test_parse_order_snapshot(
2793 #[case] price: &str,
2794 #[case] trigger: &str,
2795 #[case] expected_price: Option<Price>,
2796 #[case] expected_trigger: Option<Price>,
2797 ) {
2798 let instrument = linear_instrument();
2799 let json = load_test_json("ws_account_order.json");
2800 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
2801 value["data"][0]["price"] = serde_json::Value::String(price.to_string());
2802 value["data"][0]["triggerPrice"] = serde_json::Value::String(trigger.to_string());
2803 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
2804 serde_json::from_value(value).unwrap();
2805
2806 let snapshot = parse_order_snapshot(&msg.data[0], &instrument).unwrap();
2807 assert_eq!(snapshot.price, expected_price);
2808 assert_eq!(snapshot.trigger_price, expected_trigger);
2809 assert_eq!(snapshot.quantity, Quantity::from("0.010"));
2810 }
2811
2812 #[rstest]
2813 fn test_parse_order_snapshot_invalid_qty_returns_none() {
2814 let instrument = linear_instrument();
2815 let json = load_test_json("ws_account_order.json");
2816 let mut value: serde_json::Value = serde_json::from_str(&json).unwrap();
2817 value["data"][0]["qty"] = serde_json::Value::String(String::new());
2818 let msg: crate::websocket::messages::BybitWsAccountOrderMsg =
2819 serde_json::from_value(value).unwrap();
2820
2821 assert!(parse_order_snapshot(&msg.data[0], &instrument).is_none());
2822 }
2823
2824 #[rstest]
2825 fn test_dispatch_order_updated_on_partially_filled_price_change() {
2826 let mut ctx = DispatchTestContext::new();
2827 let value = new_order_value();
2828 ctx.accept_order(&value);
2829
2830 let mut amended = value;
2831 amended["data"][0]["orderStatus"] =
2832 serde_json::Value::String("PartiallyFilled".to_string());
2833 amended["data"][0]["cumExecQty"] = serde_json::Value::String("0.005".to_string());
2834 amended["data"][0]["price"] = serde_json::Value::String("31000".to_string());
2835 ctx.dispatch_value(&amended);
2836
2837 let updated = ctx.recv_updated();
2838 assert_eq!(updated.client_order_id, ClientOrderId::from("client-1"));
2839 assert_eq!(updated.price, Some(Price::from("31000.00")));
2840 }
2841}