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