Skip to main content

nautilus_bybit/websocket/
dispatch.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! WebSocket message dispatch for the Bybit execution client.
17//!
18//! Routes incoming [`BybitWsMessage`] variants to the appropriate parsing and
19//! event emission paths. Tracked orders (submitted through this client) produce
20//! proper order events; untracked orders fall back to execution reports for
21//! downstream reconciliation.
22
23use 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/// Order identity context stored at submission time, used by the WS dispatch
80/// task to produce proper order events without Cache access.
81#[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/// Tracks which type of WS request is pending for a given req_id.
91#[derive(Debug, Clone, Copy)]
92pub enum PendingOperation {
93    Place,
94    Cancel,
95    Amend,
96}
97
98/// Shared state for cross-stream event deduplication between the private
99/// and trade WebSocket dispatch loops.
100pub type PendingRequestData = (
101    Vec<ClientOrderId>,
102    Vec<Option<VenueOrderId>>,
103    PendingOperation,
104);
105
106/// Snapshot of an order's price, quantity, and trigger price at last dispatch.
107/// Used to detect modifications when Bybit sends back an order with the same
108/// status but changed fields.
109#[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
196/// Dispatches a WebSocket message with cross-stream deduplication.
197///
198/// For orders with a tracked identity (submitted through this client), produces
199/// proper order events (OrderAccepted, OrderCanceled, OrderFilled, etc.).
200/// For untracked orders (external or pre-existing), falls back to execution
201/// reports for downstream reconciliation.
202pub 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
297/// Dispatches a single order status update.
298///
299/// Tracked orders produce lifecycle events (OrderAccepted, OrderTriggered,
300/// OrderCanceled, OrderRejected). Untracked orders fall back to
301/// `OrderStatusReport` for reconciliation.
302fn 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                // BBO orders resolve their limit price venue-side: emit
362                // OrderUpdated after OrderAccepted when the seed diverges.
363                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                    // Partially filled then rejected - treat as canceled
456                    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                // Fills arrive on the execution channel; no event needed here.
499                // Ensure accepted was emitted so the fill has a valid prior state.
500                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                // A successful amend on a partially filled order keeps the
511                // PartiallyFilled status. Detect price/qty/trigger changes and
512                // emit OrderUpdated so PendingUpdate resolves.
513                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                // Fills arrive on the execution channel; no event needed here.
539                // Ensure accepted was emitted so the fill has a valid prior state.
540                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                // Reconcile seed against venue values before the fill (BBO
551                // orders may land directly on Filled).
552                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                // Identity cleaned up in dispatch_execution_fill when leaves_qty
576                // reaches zero, since there is no guaranteed ordering between
577                // the order and execution topics.
578            }
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                // Bybit reports a post-only order that would take liquidity as
590                // Cancelled with rejectReason=EC_PostOnlyWillTakeLiquidity,
591                // not Rejected. Surface it as OrderRejected carrying due_post_only.
592                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        // Untracked order: fall back to report for reconciliation
633        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
640/// Dispatches a single execution (fill) message.
641///
642/// Tracked orders are parsed directly to [`OrderFilled`]. Untracked orders
643/// fall back to [`FillReport`] for reconciliation.
644fn 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        // Untracked: fall back to FillReport for reconciliation
712        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
719/// Dispatches a single fast-execution (fill) message.
720///
721/// The fast channel lacks fee and exec-type fields, so all fills route to
722/// [`FillReport`] with zero commission; liquidity side is derived from the
723/// payload's `isMaker` flag. Subscribe to the standard `execution` channel
724/// for full fill metadata (e.g. fees).
725fn 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
763/// Parses a Bybit execution message directly into an [`OrderFilled`] event.
764fn 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
832/// Handles a Bybit WS order command response.
833fn 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        // Check for per-order failures in batch retExtInfo even on success
841        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                // Extract orderLinkId from the corresponding data entry
866                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    // Remove the pending request entry (if any) to get client_order_ids and op
906    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    // Single-order rejection path
985    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/// Emits the appropriate rejection event based on the pending operation type.
1033#[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
1080/// Maps an operation string to a `PendingOperation`.
1081fn 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
1097/// Parses an order snapshot from a WS order message for modification detection.
1098fn 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
1129/// Returns whether the incoming snapshot differs from the stored snapshot.
1130fn 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
1155/// Synthesizes and emits `OrderAccepted` if one has not yet been emitted for
1156/// this order. Handles fast-filling orders that skip the `New` state on Bybit.
1157fn 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
1185/// Removes a terminal order from all tracking sets.
1186fn 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
1195/// Tries to extract `orderLinkId` from the response data Value.
1196fn 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
1203/// Tries to extract `orderId` from the response data Value.
1204fn 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
1239/// Enqueues an auto-repay for a fully-filled SPOT BUY.
1240fn 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        // Fixture has orderStatus=Cancelled
1449        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        // First: synthesized Accepted
1471        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        // Second: Canceled (from Cancelled status)
1478        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        // Bybit reports a post-only order that would take liquidity as
1496        // orderStatus=Cancelled with rejectReason=EC_PostOnlyWillTakeLiquidity.
1497        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        // No identity registered → untracked
1540        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        // Register identity for the execution's orderLinkId
1570        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        // First event should be synthesized Accepted
1588        let event1 = rx.try_recv().unwrap();
1589        assert!(
1590            matches!(event1, ExecutionEvent::Order(OrderEventAny::Accepted(_))),
1591            "Expected Accepted, found {event1:?}"
1592        );
1593
1594        // Second event should be OrderFilled
1595        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        // Taker fast fill (orderLinkId populated) so the identity lookup hits.
1709        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        // First event: synthesized OrderAccepted from the identity hit.
1730        let event1 = rx.try_recv().unwrap();
1731        assert!(
1732            matches!(event1, ExecutionEvent::Order(OrderEventAny::Accepted(_))),
1733            "Expected Accepted, found {event1:?}",
1734        );
1735
1736        // Second event: FillReport carrying the hedge venue_position_id.
1737        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        // Maker fast fill: orderLinkId is empty per venue docs, so no identity lookup.
1756        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        // No OrderAccepted is synthesized for untracked fills.
1768        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        // No identity registered → untracked
1793        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        // Fixture has orderStatus=Cancelled. Patch to New for this dedup test.
1909        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        // Dispatch the same message again: dedup should suppress the duplicate
1939        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}