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