Skip to main content

nautilus_deribit/websocket/
handler.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 handler for Deribit.
17//!
18//! The handler runs in a dedicated Tokio task as the I/O boundary between the client
19//! orchestrator and the network layer. It exclusively owns the `WebSocketClient` and
20//! processes commands from the client via an unbounded channel.
21
22use std::{
23    collections::VecDeque,
24    sync::{
25        Arc,
26        atomic::{AtomicBool, AtomicU64, Ordering},
27    },
28};
29
30use ahash::AHashMap;
31use nautilus_common::cache::fifo::FifoCacheMap;
32use nautilus_core::{AtomicSet, AtomicTime, UUID4, UnixNanos, time::get_atomic_clock_realtime};
33use nautilus_model::{
34    data::{Bar, CustomData, Data, DataType, InstrumentStatus},
35    enums::{MarketStatusAction, OrderSide, OrderType},
36    events::{
37        AccountState, OrderAccepted, OrderCancelRejected, OrderFilled, OrderModifyRejected,
38        OrderRejected,
39    },
40    identifiers::{
41        AccountId, ClientOrderId, InstrumentId, StrategyId, Symbol, TraderId, VenueOrderId,
42    },
43    instruments::{Instrument, InstrumentAny},
44};
45use nautilus_network::{
46    RECONNECTED,
47    retry::{RetryManager, create_websocket_retry_manager},
48    websocket::{AuthTracker, SubscriptionState, WebSocketClient},
49};
50use parking_lot::Mutex;
51use rust_decimal::Decimal;
52use tokio_tungstenite::tungstenite::Message;
53use ustr::Ustr;
54
55use super::{
56    enums::{DeribitBookMsgType, DeribitHeartbeatType, DeribitWsChannel, DeribitWsMethod},
57    error::DeribitWsError,
58    messages::{
59        DeribitAuthResult, DeribitBookMsg, DeribitCancelAllByInstrumentParams, DeribitCancelParams,
60        DeribitChartMsg, DeribitEditParams, DeribitHeartbeatParams, DeribitInstrumentStateMsg,
61        DeribitJsonRpcRequest, DeribitOrderMsg, DeribitOrderParams, DeribitOrderResponse,
62        DeribitPerpetualMsg, DeribitPortfolioMsg, DeribitQuoteMsg, DeribitSubscribeParams,
63        DeribitTickerMsg, DeribitTradeMsg, DeribitUserTradeMsg, DeribitVolatilityIndexMsg,
64        DeribitWsMessage, NautilusWsMessage, parse_raw_message,
65    },
66    parse::{
67        OrderEventType, determine_order_event_type, parse_book_msg, parse_chart_msg,
68        parse_deribit_order_type, parse_order_accepted_with_client_order_id,
69        parse_order_canceled_with_client_order_id, parse_order_expired_with_client_order_id,
70        parse_order_updated_with_client_order_id, parse_perpetual_to_funding_rate, parse_quote_msg,
71        parse_ticker_to_index_price, parse_ticker_to_mark_price, parse_ticker_to_option_greeks,
72        parse_trades_data, parse_user_order_msg, parse_user_trade_msg, resolution_to_bar_type,
73    },
74};
75use crate::{
76    common::{
77        consts::{DERIBIT_POST_ONLY_ERROR_CODE, DERIBIT_RATE_LIMIT_KEY_ORDER, DERIBIT_VENUE},
78        enums::DeribitInstrumentState,
79        parse::{parse_portfolio_to_account_state, use_cost_for_bar_volume},
80    },
81    data_types::DeribitVolatilityIndex,
82};
83
84/// Type of pending request for request ID correlation.
85#[derive(Debug, Clone)]
86pub enum PendingRequestType {
87    /// Authentication request.
88    Authenticate,
89    /// Subscribe request with requested channels.
90    Subscribe { channels: Vec<String> },
91    /// Unsubscribe request with requested channels.
92    Unsubscribe { channels: Vec<String> },
93    /// Set heartbeat request.
94    SetHeartbeat,
95    /// Test/ping request (heartbeat response).
96    Test,
97    /// Buy order request.
98    Buy {
99        client_order_id: ClientOrderId,
100        trader_id: TraderId,
101        strategy_id: StrategyId,
102        instrument_id: InstrumentId,
103        order_side: OrderSide,
104        order_type: OrderType,
105    },
106    /// Sell order request.
107    Sell {
108        client_order_id: ClientOrderId,
109        trader_id: TraderId,
110        strategy_id: StrategyId,
111        instrument_id: InstrumentId,
112        order_side: OrderSide,
113        order_type: OrderType,
114    },
115    /// Edit order request.
116    Edit {
117        client_order_id: ClientOrderId,
118        trader_id: TraderId,
119        strategy_id: StrategyId,
120        instrument_id: InstrumentId,
121    },
122    /// Cancel order request.
123    Cancel {
124        client_order_id: ClientOrderId,
125        trader_id: TraderId,
126        strategy_id: StrategyId,
127        instrument_id: InstrumentId,
128    },
129    /// Cancel all orders by instrument request.
130    CancelAllByInstrument { instrument_id: InstrumentId },
131    /// Get order state request.
132    GetOrderState {
133        client_order_id: ClientOrderId,
134        trader_id: TraderId,
135        strategy_id: StrategyId,
136        instrument_id: InstrumentId,
137    },
138}
139
140/// Commands sent from the client to the handler.
141#[allow(missing_debug_implementations)]
142pub enum HandlerCommand {
143    /// Set the active WebSocket client.
144    SetClient(WebSocketClient),
145    /// Disconnect the WebSocket.
146    Disconnect,
147    /// Authenticate with credentials.
148    Authenticate {
149        /// Serialized auth params (DeribitAuthParams or DeribitRefreshTokenParams).
150        auth_params: serde_json::Value,
151    },
152    /// Enable heartbeat with interval.
153    SetHeartbeat { interval: u64 },
154    /// Initialize the instrument cache.
155    InitializeInstruments(Vec<InstrumentAny>),
156    /// Update a single instrument in the cache.
157    UpdateInstrument(Box<InstrumentAny>),
158    /// Subscribe to channels.
159    Subscribe { channels: Vec<String> },
160    /// Unsubscribe from channels.
161    Unsubscribe { channels: Vec<String> },
162    /// Submit a buy order.
163    Buy {
164        params: DeribitOrderParams,
165        client_order_id: ClientOrderId,
166        trader_id: TraderId,
167        strategy_id: StrategyId,
168        instrument_id: InstrumentId,
169    },
170    /// Submit a sell order.
171    Sell {
172        params: DeribitOrderParams,
173        client_order_id: ClientOrderId,
174        trader_id: TraderId,
175        strategy_id: StrategyId,
176        instrument_id: InstrumentId,
177    },
178    /// Edit an existing order.
179    Edit {
180        params: DeribitEditParams,
181        client_order_id: ClientOrderId,
182        trader_id: TraderId,
183        strategy_id: StrategyId,
184        instrument_id: InstrumentId,
185    },
186    /// Cancel an existing order.
187    Cancel {
188        params: DeribitCancelParams,
189        client_order_id: ClientOrderId,
190        trader_id: TraderId,
191        strategy_id: StrategyId,
192        instrument_id: InstrumentId,
193    },
194    /// Cancel all orders by instrument.
195    CancelAllByInstrument {
196        params: DeribitCancelAllByInstrumentParams,
197        instrument_id: InstrumentId,
198    },
199    /// Get order state.
200    GetOrderState {
201        order_id: String,
202        client_order_id: ClientOrderId,
203        trader_id: TraderId,
204        strategy_id: StrategyId,
205        instrument_id: InstrumentId,
206    },
207}
208
209/// Context for an order submitted via this handler.
210///
211/// Stores the submitted order identity for routing live order and trade updates.
212#[derive(Debug, Clone)]
213pub struct OrderContext {
214    pub client_order_id: ClientOrderId,
215    pub trader_id: TraderId,
216    pub strategy_id: StrategyId,
217    pub instrument_id: InstrumentId,
218    pub order_side: OrderSide,
219    pub order_type: OrderType,
220    pub accepted: bool,
221    pub last_order_signature: Option<OrderSignature>,
222}
223
224/// Order fields used to identify a venue amendment.
225pub type OrderSignature = (Decimal, Option<Decimal>, Option<Decimal>);
226
227/// Deribit WebSocket feed handler.
228///
229/// Runs in a dedicated Tokio task, processing commands and raw WebSocket messages.
230#[allow(missing_debug_implementations)]
231pub struct DeribitWsFeedHandler {
232    clock: &'static AtomicTime,
233    signal: Arc<AtomicBool>,
234    inner: Option<WebSocketClient>,
235    cmd_rx: tokio::sync::mpsc::UnboundedReceiver<HandlerCommand>,
236    raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
237    out_tx: tokio::sync::mpsc::UnboundedSender<NautilusWsMessage>,
238    auth_tracker: AuthTracker,
239    subscriptions_state: SubscriptionState,
240    retry_manager: RetryManager<DeribitWsError>,
241    instruments_cache: AHashMap<Ustr, InstrumentAny>,
242    option_greeks_subs: Arc<AtomicSet<InstrumentId>>,
243    mark_price_subs: Arc<AtomicSet<InstrumentId>>,
244    index_price_subs: Arc<AtomicSet<InstrumentId>>,
245    request_id_counter: AtomicU64,
246    pending_requests: AHashMap<u64, PendingRequestType>,
247    account_id: Option<AccountId>,
248    order_contexts: AHashMap<VenueOrderId, OrderContext>,
249    submitted_order_contexts: FifoCacheMap<ClientOrderId, OrderContext, 10_000>,
250
251    // Retain recent terminal identities so late trade frames stay on the order-event path
252    terminal_order_contexts: FifoCacheMap<VenueOrderId, OrderContext, 10_000>,
253    pending_bars: AHashMap<String, Bar>,
254    bars_timestamp_on_close: bool,
255    last_account_states: AHashMap<String, AccountState>,
256    book_sequence: AHashMap<Ustr, u64>,
257    pending_book_resync: Vec<String>,
258    pending_outgoing: VecDeque<NautilusWsMessage>,
259    subscribe_errors: Arc<Mutex<Vec<String>>>,
260}
261
262impl DeribitWsFeedHandler {
263    /// Creates a new feed handler.
264    #[expect(clippy::too_many_arguments)]
265    #[must_use]
266    pub fn new(
267        signal: Arc<AtomicBool>,
268        cmd_rx: tokio::sync::mpsc::UnboundedReceiver<HandlerCommand>,
269        raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
270        out_tx: tokio::sync::mpsc::UnboundedSender<NautilusWsMessage>,
271        auth_tracker: AuthTracker,
272        subscriptions_state: SubscriptionState,
273        option_greeks_subs: Arc<AtomicSet<InstrumentId>>,
274        mark_price_subs: Arc<AtomicSet<InstrumentId>>,
275        index_price_subs: Arc<AtomicSet<InstrumentId>>,
276        account_id: Option<AccountId>,
277        bars_timestamp_on_close: bool,
278        subscribe_errors: Arc<Mutex<Vec<String>>>,
279    ) -> Self {
280        Self {
281            clock: get_atomic_clock_realtime(),
282            signal,
283            inner: None,
284            cmd_rx,
285            raw_rx,
286            out_tx,
287            auth_tracker,
288            subscriptions_state,
289            retry_manager: create_websocket_retry_manager(),
290            instruments_cache: AHashMap::new(),
291            option_greeks_subs,
292            mark_price_subs,
293            index_price_subs,
294            request_id_counter: AtomicU64::new(1),
295            pending_requests: AHashMap::new(),
296            account_id,
297            order_contexts: AHashMap::new(),
298            submitted_order_contexts: FifoCacheMap::new(),
299            terminal_order_contexts: FifoCacheMap::new(),
300            pending_bars: AHashMap::new(),
301            bars_timestamp_on_close,
302            last_account_states: AHashMap::new(),
303            book_sequence: AHashMap::new(),
304            pending_book_resync: Vec::new(),
305            pending_outgoing: VecDeque::new(),
306            subscribe_errors,
307        }
308    }
309
310    /// Sets the account ID for order/fill reports.
311    pub fn set_account_id(&mut self, account_id: AccountId) {
312        self.account_id = Some(account_id);
313    }
314
315    /// Returns the account ID.
316    #[must_use]
317    pub fn account_id(&self) -> Option<AccountId> {
318        self.account_id
319    }
320
321    fn clear_state(&mut self) {
322        let pending_count = self.pending_requests.len();
323        let bars_count = self.pending_bars.len();
324        let account_count = self.last_account_states.len();
325        let book_count = self.book_sequence.len();
326        let outgoing_count = self.pending_outgoing.len();
327
328        self.pending_requests.clear();
329        self.pending_bars.clear();
330        self.last_account_states.clear();
331        self.book_sequence.clear();
332        self.pending_book_resync.clear();
333        self.pending_outgoing.clear();
334
335        log::debug!(
336            "Reset state: pending_requests={pending_count}, pending_bars={bars_count}, \
337            account_states={account_count}, book_sequence={book_count}, \
338            pending_outgoing={outgoing_count}"
339        );
340    }
341
342    /// Generates a unique request ID.
343    fn next_request_id(&self) -> u64 {
344        self.request_id_counter.fetch_add(1, Ordering::Relaxed)
345    }
346
347    /// Returns the current timestamp.
348    fn ts_init(&self) -> UnixNanos {
349        self.clock.get_time_ns()
350    }
351
352    async fn send_tracked_request(
353        &mut self,
354        request_id: u64,
355        payload: Result<String, DeribitWsError>,
356        rate_limit_keys: Option<&[Ustr]>,
357    ) -> Result<(), DeribitWsError> {
358        let payload = match payload {
359            Ok(p) => p,
360            Err(e) => {
361                self.pending_requests.remove(&request_id);
362                return Err(e);
363            }
364        };
365        self.send_with_retry(payload, rate_limit_keys).await
366    }
367
368    /// Sends a message over the WebSocket with retry logic.
369    async fn send_with_retry(
370        &self,
371        payload: String,
372        rate_limit_keys: Option<&[Ustr]>,
373    ) -> Result<(), DeribitWsError> {
374        if let Some(client) = &self.inner {
375            let keys_owned: Option<Vec<Ustr>> = rate_limit_keys.map(|k| k.to_vec());
376            self.retry_manager
377                .execute_with_retry(
378                    "websocket_send",
379                    || {
380                        let payload = payload.clone();
381                        let keys = keys_owned.clone();
382                        async move {
383                            client
384                                .send_text(payload, keys.as_deref())
385                                .await
386                                .map_err(|e| DeribitWsError::Send(e.to_string()))
387                        }
388                    },
389                    |e| matches!(e, DeribitWsError::Send(_)),
390                    |e| DeribitWsError::Timeout(e.to_string()),
391                )
392                .await
393        } else {
394            Err(DeribitWsError::NotConnected)
395        }
396    }
397
398    /// Handles a subscribe command.
399    ///
400    /// Note: The client has already called `mark_subscribe` before sending this command.
401    async fn handle_subscribe(&mut self, channels: Vec<String>) -> Result<(), DeribitWsError> {
402        let request_id = self.next_request_id();
403
404        // Track this request for response correlation
405        self.pending_requests.insert(
406            request_id,
407            PendingRequestType::Subscribe {
408                channels: channels.clone(),
409            },
410        );
411
412        // Deribit requires private/subscribe for authenticated channels
413        let method = if channels
414            .iter()
415            .any(|ch| DeribitWsChannel::requires_auth(ch))
416        {
417            DeribitWsMethod::PrivateSubscribe
418        } else {
419            DeribitWsMethod::PublicSubscribe
420        };
421
422        let request = DeribitJsonRpcRequest::new(
423            request_id,
424            method.as_method_str(),
425            DeribitSubscribeParams {
426                channels: channels.clone(),
427            },
428        );
429
430        let payload =
431            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
432
433        log::debug!("Subscribing to channels: request_id={request_id}, channels={channels:?}");
434        self.send_tracked_request(request_id, payload, None).await
435    }
436
437    /// Handles an unsubscribe command.
438    async fn handle_unsubscribe(&mut self, channels: Vec<String>) -> Result<(), DeribitWsError> {
439        let request_id = self.next_request_id();
440
441        // Track this request for response correlation
442        self.pending_requests.insert(
443            request_id,
444            PendingRequestType::Unsubscribe {
445                channels: channels.clone(),
446            },
447        );
448
449        let method = if channels
450            .iter()
451            .any(|ch| DeribitWsChannel::requires_auth(ch))
452        {
453            DeribitWsMethod::PrivateUnsubscribe
454        } else {
455            DeribitWsMethod::PublicUnsubscribe
456        };
457
458        let request = DeribitJsonRpcRequest::new(
459            request_id,
460            method.as_method_str(),
461            DeribitSubscribeParams {
462                channels: channels.clone(),
463            },
464        );
465
466        let payload =
467            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
468
469        log::debug!("Unsubscribing from channels: request_id={request_id}, channels={channels:?}");
470        self.send_tracked_request(request_id, payload, None).await
471    }
472
473    /// Handles enabling heartbeat.
474    async fn handle_set_heartbeat(&mut self, interval: u64) -> Result<(), DeribitWsError> {
475        let request_id = self.next_request_id();
476
477        // Track this request for response correlation
478        self.pending_requests
479            .insert(request_id, PendingRequestType::SetHeartbeat);
480
481        let request = DeribitJsonRpcRequest::new(
482            request_id,
483            DeribitWsMethod::SetHeartbeat.as_method_str(),
484            DeribitHeartbeatParams { interval },
485        );
486
487        let payload =
488            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
489
490        log::debug!(
491            "Enabling heartbeat with interval: request_id={request_id}, interval={interval} seconds"
492        );
493        self.send_tracked_request(request_id, payload, None).await
494    }
495
496    /// Responds to a heartbeat test_request.
497    async fn handle_heartbeat_test_request(&mut self) -> Result<(), DeribitWsError> {
498        let request_id = self.next_request_id();
499
500        // Track this request for response correlation
501        self.pending_requests
502            .insert(request_id, PendingRequestType::Test);
503
504        let request = DeribitJsonRpcRequest::new(
505            request_id,
506            DeribitWsMethod::Test.as_method_str(),
507            serde_json::json!({}),
508        );
509
510        let payload =
511            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
512
513        log::trace!("Responding to heartbeat test_request: request_id={request_id}");
514        self.send_tracked_request(request_id, payload, None).await
515    }
516
517    /// Handles a buy order command.
518    async fn handle_buy(
519        &mut self,
520        params: DeribitOrderParams,
521        client_order_id: ClientOrderId,
522        trader_id: TraderId,
523        strategy_id: StrategyId,
524        instrument_id: InstrumentId,
525    ) -> Result<(), DeribitWsError> {
526        let request_id = self.next_request_id();
527        let order_type = parse_deribit_order_type(&params.order_type);
528        let order_signature = (params.amount, params.price, params.trigger_price);
529
530        self.submitted_order_contexts.insert(
531            client_order_id,
532            OrderContext {
533                client_order_id,
534                trader_id,
535                strategy_id,
536                instrument_id,
537                order_side: OrderSide::Buy,
538                order_type,
539                accepted: false,
540                last_order_signature: Some(order_signature),
541            },
542        );
543
544        self.pending_requests.insert(
545            request_id,
546            PendingRequestType::Buy {
547                client_order_id,
548                trader_id,
549                strategy_id,
550                instrument_id,
551                order_side: OrderSide::Buy,
552                order_type,
553            },
554        );
555
556        let request =
557            DeribitJsonRpcRequest::new(request_id, DeribitWsMethod::Buy.as_method_str(), params);
558
559        let payload =
560            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
561
562        log::debug!("Sending buy order: request_id={request_id}");
563        self.send_tracked_request(
564            request_id,
565            payload,
566            Some(DERIBIT_RATE_LIMIT_KEY_ORDER.as_slice()),
567        )
568        .await
569    }
570
571    /// Handles a sell order command.
572    async fn handle_sell(
573        &mut self,
574        params: DeribitOrderParams,
575        client_order_id: ClientOrderId,
576        trader_id: TraderId,
577        strategy_id: StrategyId,
578        instrument_id: InstrumentId,
579    ) -> Result<(), DeribitWsError> {
580        let request_id = self.next_request_id();
581        let order_type = parse_deribit_order_type(&params.order_type);
582        let order_signature = (params.amount, params.price, params.trigger_price);
583
584        self.submitted_order_contexts.insert(
585            client_order_id,
586            OrderContext {
587                client_order_id,
588                trader_id,
589                strategy_id,
590                instrument_id,
591                order_side: OrderSide::Sell,
592                order_type,
593                accepted: false,
594                last_order_signature: Some(order_signature),
595            },
596        );
597
598        self.pending_requests.insert(
599            request_id,
600            PendingRequestType::Sell {
601                client_order_id,
602                trader_id,
603                strategy_id,
604                instrument_id,
605                order_side: OrderSide::Sell,
606                order_type,
607            },
608        );
609
610        let request =
611            DeribitJsonRpcRequest::new(request_id, DeribitWsMethod::Sell.as_method_str(), params);
612
613        let payload =
614            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
615
616        log::debug!("Sending sell order: request_id={request_id}");
617        self.send_tracked_request(
618            request_id,
619            payload,
620            Some(DERIBIT_RATE_LIMIT_KEY_ORDER.as_slice()),
621        )
622        .await
623    }
624
625    /// Handles an edit order command.
626    async fn handle_edit(
627        &mut self,
628        params: DeribitEditParams,
629        client_order_id: ClientOrderId,
630        trader_id: TraderId,
631        strategy_id: StrategyId,
632        instrument_id: InstrumentId,
633    ) -> Result<(), DeribitWsError> {
634        let request_id = self.next_request_id();
635        let order_id = params.order_id.clone();
636
637        self.pending_requests.insert(
638            request_id,
639            PendingRequestType::Edit {
640                client_order_id,
641                trader_id,
642                strategy_id,
643                instrument_id,
644            },
645        );
646
647        let request =
648            DeribitJsonRpcRequest::new(request_id, DeribitWsMethod::Edit.as_method_str(), params);
649
650        let payload =
651            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
652
653        log::debug!("Sending edit order: request_id={request_id}, order_id={order_id}");
654        self.send_tracked_request(
655            request_id,
656            payload,
657            Some(DERIBIT_RATE_LIMIT_KEY_ORDER.as_slice()),
658        )
659        .await
660    }
661
662    /// Handles a cancel order command.
663    async fn handle_cancel(
664        &mut self,
665        params: DeribitCancelParams,
666        client_order_id: ClientOrderId,
667        trader_id: TraderId,
668        strategy_id: StrategyId,
669        instrument_id: InstrumentId,
670    ) -> Result<(), DeribitWsError> {
671        let request_id = self.next_request_id();
672        let order_id = params.order_id.clone();
673
674        self.pending_requests.insert(
675            request_id,
676            PendingRequestType::Cancel {
677                client_order_id,
678                trader_id,
679                strategy_id,
680                instrument_id,
681            },
682        );
683
684        let request =
685            DeribitJsonRpcRequest::new(request_id, DeribitWsMethod::Cancel.as_method_str(), params);
686
687        let payload =
688            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
689
690        log::debug!("Sending cancel order: request_id={request_id}, order_id={order_id}");
691        self.send_tracked_request(
692            request_id,
693            payload,
694            Some(DERIBIT_RATE_LIMIT_KEY_ORDER.as_slice()),
695        )
696        .await
697    }
698
699    /// Handles cancel all orders by instrument command.
700    async fn handle_cancel_all_by_instrument(
701        &mut self,
702        params: DeribitCancelAllByInstrumentParams,
703        instrument_id: InstrumentId,
704    ) -> Result<(), DeribitWsError> {
705        let request_id = self.next_request_id();
706        let instrument_name = params.instrument_name.clone();
707
708        // Track this request for response correlation
709        self.pending_requests.insert(
710            request_id,
711            PendingRequestType::CancelAllByInstrument { instrument_id },
712        );
713
714        let request = DeribitJsonRpcRequest::new(
715            request_id,
716            DeribitWsMethod::CancelAllByInstrument.as_method_str(),
717            params,
718        );
719
720        let payload =
721            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
722
723        log::debug!(
724            "Sending cancel_all_by_instrument: request_id={request_id}, instrument={instrument_name}"
725        );
726        self.send_tracked_request(
727            request_id,
728            payload,
729            Some(DERIBIT_RATE_LIMIT_KEY_ORDER.as_slice()),
730        )
731        .await
732    }
733
734    /// Handles get order state command.
735    async fn handle_get_order_state(
736        &mut self,
737        order_id: String,
738        client_order_id: ClientOrderId,
739        trader_id: TraderId,
740        strategy_id: StrategyId,
741        instrument_id: InstrumentId,
742    ) -> Result<(), DeribitWsError> {
743        let request_id = self.next_request_id();
744
745        // Track this request for response correlation
746        self.pending_requests.insert(
747            request_id,
748            PendingRequestType::GetOrderState {
749                client_order_id,
750                trader_id,
751                strategy_id,
752                instrument_id,
753            },
754        );
755
756        let params = serde_json::json!({
757            "order_id": order_id
758        });
759
760        let request = DeribitJsonRpcRequest::new(
761            request_id,
762            DeribitWsMethod::GetOrderState.as_method_str(),
763            params,
764        );
765
766        let payload =
767            serde_json::to_string(&request).map_err(|e| DeribitWsError::Json(e.to_string()));
768
769        log::debug!("Sending get_order_state: request_id={request_id}, order_id={order_id}");
770        self.send_tracked_request(
771            request_id,
772            payload,
773            Some(DERIBIT_RATE_LIMIT_KEY_ORDER.as_slice()),
774        )
775        .await
776    }
777
778    /// Processes a command from the client.
779    async fn process_command(&mut self, cmd: HandlerCommand) {
780        match cmd {
781            HandlerCommand::SetClient(client) => {
782                log::debug!("Setting WebSocket client");
783                self.inner = Some(client);
784            }
785            HandlerCommand::Disconnect => {
786                log::debug!("Disconnecting WebSocket");
787
788                if let Some(client) = self.inner.take() {
789                    client.disconnect().await;
790                }
791            }
792            HandlerCommand::Authenticate { auth_params } => {
793                let request_id = self.next_request_id();
794                log::debug!("Authenticating: request_id={request_id}");
795
796                // Track this request for response correlation
797                self.pending_requests
798                    .insert(request_id, PendingRequestType::Authenticate);
799
800                let request = DeribitJsonRpcRequest::new(
801                    request_id,
802                    DeribitWsMethod::PublicAuth.as_method_str(),
803                    auth_params,
804                );
805
806                match serde_json::to_string(&request) {
807                    Ok(payload) => {
808                        if let Err(e) = self.send_with_retry(payload, None).await {
809                            self.pending_requests.remove(&request_id);
810                            log::error!("Authentication send failed: {e}");
811                            self.auth_tracker.fail(format!("Send failed: {e}"));
812                        }
813                    }
814                    Err(e) => {
815                        self.pending_requests.remove(&request_id);
816                        log::error!("Failed to serialize auth request: {e}");
817                        self.auth_tracker.fail(format!("Serialization failed: {e}"));
818                    }
819                }
820            }
821            HandlerCommand::SetHeartbeat { interval } => {
822                if let Err(e) = self.handle_set_heartbeat(interval).await {
823                    log::error!("Set heartbeat failed: {e}");
824                }
825            }
826            HandlerCommand::InitializeInstruments(instruments) => {
827                log::debug!("Handler received {} instruments", instruments.len());
828                self.instruments_cache.clear();
829                for inst in instruments {
830                    self.instruments_cache
831                        .insert(inst.raw_symbol().inner(), inst);
832                }
833            }
834            HandlerCommand::UpdateInstrument(instrument) => {
835                log::trace!("Updating instrument: {}", instrument.raw_symbol());
836                self.instruments_cache
837                    .insert(instrument.raw_symbol().inner(), *instrument);
838            }
839            HandlerCommand::Subscribe { channels } => {
840                if let Err(e) = self.handle_subscribe(channels).await {
841                    log::error!("Subscribe failed: {e}");
842                }
843            }
844            HandlerCommand::Unsubscribe { channels } => {
845                // User-initiated unsubscribe cancels any pending book resync
846                // for these channels so we don't re-subscribe against user intent
847                self.pending_book_resync.retain(|ch| !channels.contains(ch));
848
849                if let Err(e) = self.handle_unsubscribe(channels).await {
850                    log::error!("Unsubscribe failed: {e}");
851                }
852            }
853            HandlerCommand::Buy {
854                params,
855                client_order_id,
856                trader_id,
857                strategy_id,
858                instrument_id,
859            } => {
860                if let Err(e) = self
861                    .handle_buy(
862                        params,
863                        client_order_id,
864                        trader_id,
865                        strategy_id,
866                        instrument_id,
867                    )
868                    .await
869                {
870                    log::error!("Buy order failed: {e}");
871                }
872            }
873            HandlerCommand::Sell {
874                params,
875                client_order_id,
876                trader_id,
877                strategy_id,
878                instrument_id,
879            } => {
880                if let Err(e) = self
881                    .handle_sell(
882                        params,
883                        client_order_id,
884                        trader_id,
885                        strategy_id,
886                        instrument_id,
887                    )
888                    .await
889                {
890                    log::error!("Sell order failed: {e}");
891                }
892            }
893            HandlerCommand::Edit {
894                params,
895                client_order_id,
896                trader_id,
897                strategy_id,
898                instrument_id,
899            } => {
900                if let Err(e) = self
901                    .handle_edit(
902                        params,
903                        client_order_id,
904                        trader_id,
905                        strategy_id,
906                        instrument_id,
907                    )
908                    .await
909                {
910                    log::error!("Edit order failed: {e}");
911                }
912            }
913            HandlerCommand::Cancel {
914                params,
915                client_order_id,
916                trader_id,
917                strategy_id,
918                instrument_id,
919            } => {
920                if let Err(e) = self
921                    .handle_cancel(
922                        params,
923                        client_order_id,
924                        trader_id,
925                        strategy_id,
926                        instrument_id,
927                    )
928                    .await
929                {
930                    log::error!("Cancel order failed: {e}");
931                }
932            }
933            HandlerCommand::CancelAllByInstrument {
934                params,
935                instrument_id,
936            } => {
937                if let Err(e) = self
938                    .handle_cancel_all_by_instrument(params, instrument_id)
939                    .await
940                {
941                    log::error!("Cancel all by instrument failed: {e}");
942                }
943            }
944            HandlerCommand::GetOrderState {
945                order_id,
946                client_order_id,
947                trader_id,
948                strategy_id,
949                instrument_id,
950            } => {
951                if let Err(e) = self
952                    .handle_get_order_state(
953                        order_id,
954                        client_order_id,
955                        trader_id,
956                        strategy_id,
957                        instrument_id,
958                    )
959                    .await
960                {
961                    log::error!("Get order state failed: {e}");
962                }
963            }
964        }
965    }
966
967    /// Processes a raw WebSocket message.
968    async fn process_raw_message(&mut self, text: &str) -> Option<NautilusWsMessage> {
969        if text == RECONNECTED {
970            log::info!("Received reconnection signal");
971
972            self.auth_tracker.invalidate();
973            self.clear_state();
974
975            return Some(NautilusWsMessage::Reconnected);
976        }
977
978        // Parse the JSON-RPC message
979        let ws_msg = match parse_raw_message(text) {
980            Ok(msg) => msg,
981            Err(e) => {
982                log::warn!("Failed to parse message: {e}");
983                return None;
984            }
985        };
986
987        let ts_init = self.ts_init();
988
989        match ws_msg {
990            DeribitWsMessage::Response(response) => {
991                // Look up the request type by ID for explicit correlation
992                if let Some(request_id) = response.id
993                    && let Some(request_type) = self.pending_requests.remove(&request_id)
994                {
995                    match request_type {
996                        PendingRequestType::Authenticate => {
997                            if let Some(error) = &response.error {
998                                let reason = format!(
999                                    "Authentication error code={}: {}",
1000                                    error.code, error.message
1001                                );
1002                                log::error!(
1003                                    "Authentication failed: code={}, message={}, request_id={}",
1004                                    error.code,
1005                                    error.message,
1006                                    request_id
1007                                );
1008                                self.auth_tracker.fail(reason.clone());
1009                                return Some(NautilusWsMessage::AuthenticationFailed(reason));
1010                            } else if let Some(result) = &response.result {
1011                                match serde_json::from_value::<DeribitAuthResult>(result.clone()) {
1012                                    Ok(auth_result) => {
1013                                        self.auth_tracker.succeed();
1014                                        log::debug!(
1015                                            "WebSocket authenticated successfully (request_id={}, scope={}, expires_in={}s)",
1016                                            request_id,
1017                                            auth_result.scope,
1018                                            auth_result.expires_in
1019                                        );
1020                                        return Some(NautilusWsMessage::Authenticated(Box::new(
1021                                            auth_result,
1022                                        )));
1023                                    }
1024                                    Err(e) => {
1025                                        let reason = format!("Failed to parse auth result: {e}");
1026                                        log::error!("{reason}: request_id={request_id}");
1027                                        self.auth_tracker.fail(reason.clone());
1028                                        return Some(NautilusWsMessage::AuthenticationFailed(
1029                                            reason,
1030                                        ));
1031                                    }
1032                                }
1033                            }
1034                        }
1035                        PendingRequestType::Subscribe { channels } => {
1036                            if let Some(error) = &response.error {
1037                                log::error!(
1038                                    "Subscribe failed: code={}, message={}, channels={:?}, request_id={}",
1039                                    error.code,
1040                                    error.message,
1041                                    channels,
1042                                    request_id
1043                                );
1044
1045                                self.subscribe_errors.lock().push(format!(
1046                                    "Subscribe rejected: code={}, message={}",
1047                                    error.code, error.message,
1048                                ));
1049                            } else {
1050                                // Confirm each channel in the subscription
1051                                for ch in &channels {
1052                                    self.subscriptions_state.confirm_subscribe(ch);
1053                                    log::debug!("Subscription confirmed: {ch}");
1054                                }
1055                            }
1056                        }
1057                        PendingRequestType::Unsubscribe { channels } => {
1058                            if let Some(error) = &response.error {
1059                                log::error!(
1060                                    "Unsubscribe failed: code={}, message={}, channels={:?}, request_id={}",
1061                                    error.code,
1062                                    error.message,
1063                                    channels,
1064                                    request_id
1065                                );
1066                            } else {
1067                                for ch in &channels {
1068                                    self.subscriptions_state.confirm_unsubscribe(ch);
1069                                    log::debug!("Unsubscription confirmed: {ch}");
1070                                }
1071                            }
1072
1073                            // Resubscribe channels pending book resync (kept in
1074                            // pending_book_resync until a fresh snapshot arrives)
1075                            if !self.pending_book_resync.is_empty() {
1076                                let resync: Vec<String> = channels
1077                                    .iter()
1078                                    .filter(|ch| self.pending_book_resync.contains(ch))
1079                                    .cloned()
1080                                    .collect();
1081
1082                                if !resync.is_empty() {
1083                                    let _ = self.handle_subscribe(resync).await;
1084                                }
1085                            }
1086                        }
1087                        PendingRequestType::SetHeartbeat => {
1088                            if let Some(error) = &response.error {
1089                                log::error!(
1090                                    "Set heartbeat failed: code={}, message={}, request_id={}",
1091                                    error.code,
1092                                    error.message,
1093                                    request_id
1094                                );
1095                            } else {
1096                                log::debug!("Heartbeat enabled (request_id={request_id})");
1097                            }
1098                        }
1099                        PendingRequestType::Test => {
1100                            if let Some(error) = &response.error {
1101                                log::warn!(
1102                                    "Heartbeat test failed: code={}, message={}, request_id={}",
1103                                    error.code,
1104                                    error.message,
1105                                    request_id
1106                                );
1107                            } else {
1108                                log::trace!(
1109                                    "Heartbeat test acknowledged (request_id={request_id})"
1110                                );
1111                            }
1112                        }
1113                        PendingRequestType::Cancel {
1114                            client_order_id,
1115                            trader_id,
1116                            strategy_id,
1117                            instrument_id,
1118                        } => {
1119                            if let Some(result) = &response.result {
1120                                match serde_json::from_value::<DeribitOrderMsg>(result.clone()) {
1121                                    Ok(order_msg) => {
1122                                        let venue_order_id =
1123                                            VenueOrderId::new(order_msg.order_id.as_str());
1124                                        log::debug!(
1125                                            "Cancel confirmed: venue_order_id={venue_order_id}, \
1126                                            client_order_id={client_order_id}, state={}",
1127                                            order_msg.order_state
1128                                        );
1129
1130                                        // Emit OrderCanceled from the response path so we
1131                                        // do not lose the event during a reconnection gap.
1132                                        // Both paths check terminal context to suppress
1133                                        // duplicates regardless of which arrives first.
1134                                        if order_msg.order_state == "cancelled"
1135                                            && !self
1136                                                .terminal_order_contexts
1137                                                .contains_key(&venue_order_id)
1138                                        {
1139                                            let instrument_name_ustr = order_msg.instrument_name;
1140
1141                                            if let Some(instrument) =
1142                                                self.instruments_cache.get(&instrument_name_ustr)
1143                                                && let Some(account_id) = self.account_id
1144                                            {
1145                                                let event =
1146                                                    parse_order_canceled_with_client_order_id(
1147                                                        &order_msg,
1148                                                        instrument,
1149                                                        account_id,
1150                                                        trader_id,
1151                                                        strategy_id,
1152                                                        client_order_id,
1153                                                        ts_init,
1154                                                    );
1155                                                let context = self
1156                                                    .find_order_context(
1157                                                        venue_order_id,
1158                                                        Some(client_order_id),
1159                                                    )
1160                                                    .or_else(|| {
1161                                                        let order_side =
1162                                                            match order_msg.direction.as_str() {
1163                                                                "buy" => OrderSide::Buy,
1164                                                                "sell" => OrderSide::Sell,
1165                                                                _ => return None,
1166                                                            };
1167                                                        Some(OrderContext {
1168                                                            client_order_id,
1169                                                            trader_id,
1170                                                            strategy_id,
1171                                                            instrument_id,
1172                                                            order_side,
1173                                                            order_type: parse_deribit_order_type(
1174                                                                &order_msg.order_type,
1175                                                            ),
1176                                                            accepted: true,
1177                                                            last_order_signature: Some(
1178                                                                Self::order_signature(&order_msg),
1179                                                            ),
1180                                                        })
1181                                                    });
1182
1183                                                if let Some(context) = context {
1184                                                    self.finish_order_context(
1185                                                        venue_order_id,
1186                                                        &context,
1187                                                    );
1188                                                }
1189                                                return Some(NautilusWsMessage::OrderCanceled(
1190                                                    event,
1191                                                ));
1192                                            }
1193                                        }
1194                                    }
1195                                    Err(e) => {
1196                                        log::error!(
1197                                            "Failed to parse cancel response: request_id={request_id}, error={e}"
1198                                        );
1199                                    }
1200                                }
1201                            } else if let Some(error) = &response.error {
1202                                log::error!(
1203                                    "Cancel rejected: code={}, message={}, client_order_id={}",
1204                                    error.code,
1205                                    error.message,
1206                                    client_order_id
1207                                );
1208                                return Some(NautilusWsMessage::OrderCancelRejected(
1209                                    OrderCancelRejected::new(
1210                                        trader_id,
1211                                        strategy_id,
1212                                        instrument_id,
1213                                        client_order_id,
1214                                        ustr::ustr(&format!(
1215                                            "code={}: {}",
1216                                            error.code, error.message
1217                                        )),
1218                                        UUID4::new(),
1219                                        ts_init,
1220                                        ts_init,
1221                                        false,
1222                                        None, // venue_order_id not available in error response
1223                                        self.account_id,
1224                                    ),
1225                                ));
1226                            }
1227                        }
1228                        PendingRequestType::CancelAllByInstrument { instrument_id } => {
1229                            if let Some(result) = &response.result {
1230                                match serde_json::from_value::<u64>(result.clone()) {
1231                                    Ok(count) => {
1232                                        log::debug!(
1233                                            "Cancelled {count} orders for instrument {instrument_id}"
1234                                        );
1235                                        // Individual order status updates come via user.orders subscription
1236                                    }
1237                                    Err(e) => {
1238                                        log::warn!("Failed to parse cancel_all response: {e}");
1239                                    }
1240                                }
1241                            } else if let Some(error) = &response.error {
1242                                log::error!(
1243                                    "Cancel all by instrument rejected: code={}, message={}, instrument_id={}",
1244                                    error.code,
1245                                    error.message,
1246                                    instrument_id
1247                                );
1248                            }
1249                        }
1250                        PendingRequestType::Buy {
1251                            client_order_id,
1252                            trader_id,
1253                            strategy_id,
1254                            instrument_id,
1255                            order_side,
1256                            order_type,
1257                        }
1258                        | PendingRequestType::Sell {
1259                            client_order_id,
1260                            trader_id,
1261                            strategy_id,
1262                            instrument_id,
1263                            order_side,
1264                            order_type,
1265                        } => {
1266                            if let Some(result) = &response.result {
1267                                match serde_json::from_value::<DeribitOrderResponse>(result.clone())
1268                                {
1269                                    Ok(order_response) => {
1270                                        let venue_order_id_str = &order_response.order.order_id;
1271                                        let venue_order_id =
1272                                            VenueOrderId::new(venue_order_id_str.as_str());
1273                                        let order_state = &order_response.order.order_state;
1274                                        log::debug!(
1275                                            "Order response: venue_order_id={venue_order_id}, client_order_id={client_order_id}, state={order_state}"
1276                                        );
1277
1278                                        let mut context = self
1279                                            .find_order_context(
1280                                                venue_order_id,
1281                                                Some(client_order_id),
1282                                            )
1283                                            .unwrap_or(OrderContext {
1284                                                client_order_id,
1285                                                trader_id,
1286                                                strategy_id,
1287                                                instrument_id,
1288                                                order_side,
1289                                                order_type,
1290                                                accepted: false,
1291                                                last_order_signature: Some(Self::order_signature(
1292                                                    &order_response.order,
1293                                                )),
1294                                            });
1295                                        context.last_order_signature =
1296                                            Some(Self::order_signature(&order_response.order));
1297                                        self.bind_order_context(venue_order_id, context.clone());
1298
1299                                        if !order_response.trades.is_empty() {
1300                                            let outgoing = self
1301                                                .route_user_trades(&order_response.trades, ts_init);
1302                                            self.pending_outgoing.extend(outgoing);
1303                                        } else if order_state == "filled" {
1304                                            log::debug!(
1305                                                "Deferring acceptance for fast-filled order until its trade arrives: venue_order_id={venue_order_id}, client_order_id={client_order_id}"
1306                                            );
1307                                            self.finish_order_context(venue_order_id, &context);
1308                                        } else if context.accepted {
1309                                            log::trace!(
1310                                                "Skipping duplicate OrderAccepted response: venue_order_id={venue_order_id}"
1311                                            );
1312                                        } else {
1313                                            let instrument_name_ustr = Ustr::from(
1314                                                order_response.order.instrument_name.as_str(),
1315                                            );
1316
1317                                            if let Some(instrument) =
1318                                                self.instruments_cache.get(&instrument_name_ustr)
1319                                            {
1320                                                if let Some(account_id) = self.account_id {
1321                                                    let event =
1322                                                        parse_order_accepted_with_client_order_id(
1323                                                            &order_response.order,
1324                                                            instrument,
1325                                                            account_id,
1326                                                            trader_id,
1327                                                            strategy_id,
1328                                                            client_order_id,
1329                                                            ts_init,
1330                                                        );
1331                                                    context.accepted = true;
1332                                                    self.bind_order_context(
1333                                                        venue_order_id,
1334                                                        context,
1335                                                    );
1336                                                    return Some(NautilusWsMessage::OrderAccepted(
1337                                                        event,
1338                                                    ));
1339                                                } else {
1340                                                    log::warn!(
1341                                                        "Cannot create OrderAccepted: account_id not set"
1342                                                    );
1343                                                }
1344                                            } else {
1345                                                log::warn!(
1346                                                    "Instrument {instrument_name_ustr} not found in cache for order response"
1347                                                );
1348                                            }
1349                                        }
1350                                    }
1351                                    Err(e) => {
1352                                        log::error!(
1353                                            "Failed to parse order response: request_id={request_id}, error={e}"
1354                                        );
1355                                    }
1356                                }
1357                            } else if let Some(error) = &response.error {
1358                                let due_post_only = error.code == DERIBIT_POST_ONLY_ERROR_CODE;
1359                                let reason = if let Some(data) = &error.data {
1360                                    format!(
1361                                        "code={}: {} (data: {})",
1362                                        error.code, error.message, data
1363                                    )
1364                                } else {
1365                                    format!("code={}: {}", error.code, error.message)
1366                                };
1367
1368                                log::debug!(
1369                                    "Order rejected: {reason}, client_order_id={client_order_id}"
1370                                );
1371                                self.submitted_order_contexts.remove(&client_order_id);
1372                                return Some(NautilusWsMessage::OrderRejected(OrderRejected::new(
1373                                    trader_id,
1374                                    strategy_id,
1375                                    instrument_id,
1376                                    client_order_id,
1377                                    self.account_id.unwrap_or(AccountId::new("DERIBIT-UNKNOWN")),
1378                                    ustr::ustr(&reason),
1379                                    UUID4::new(),
1380                                    ts_init,
1381                                    ts_init,
1382                                    false,
1383                                    due_post_only,
1384                                )));
1385                            }
1386                        }
1387                        PendingRequestType::Edit {
1388                            client_order_id,
1389                            trader_id,
1390                            strategy_id,
1391                            instrument_id,
1392                        } => {
1393                            if let Some(result) = &response.result {
1394                                match serde_json::from_value::<DeribitOrderResponse>(result.clone())
1395                                {
1396                                    Ok(order_response) => {
1397                                        let venue_order_id =
1398                                            VenueOrderId::new(&order_response.order.order_id);
1399                                        log::debug!(
1400                                            "Order updated: venue_order_id={}, client_order_id={}, state={}",
1401                                            venue_order_id,
1402                                            client_order_id,
1403                                            order_response.order.order_state
1404                                        );
1405
1406                                        let instrument_name_ustr = Ustr::from(
1407                                            order_response.order.instrument_name.as_str(),
1408                                        );
1409
1410                                        if let Some(instrument) =
1411                                            self.instruments_cache.get(&instrument_name_ustr)
1412                                        {
1413                                            if let Some(account_id) = self.account_id {
1414                                                let Some(mut context) = self.find_order_context(
1415                                                    venue_order_id,
1416                                                    Some(client_order_id),
1417                                                ) else {
1418                                                    let report = parse_user_order_msg(
1419                                                        &order_response.order,
1420                                                        instrument,
1421                                                        account_id,
1422                                                        ts_init,
1423                                                    );
1424                                                    let outgoing = self.route_user_trades(
1425                                                        &order_response.trades,
1426                                                        ts_init,
1427                                                    );
1428                                                    self.pending_outgoing.extend(outgoing);
1429                                                    return report
1430                                                    .map(|report| {
1431                                                        NautilusWsMessage::OrderStatusReports(vec![
1432                                                            report,
1433                                                        ])
1434                                                    })
1435                                                    .map_err(|e| {
1436                                                        log::warn!(
1437                                                            "Failed to parse external edit response: {e}"
1438                                                        );
1439                                                    })
1440                                                    .ok();
1441                                                };
1442                                                let was_terminal = self
1443                                                    .terminal_order_contexts
1444                                                    .contains_key(&venue_order_id);
1445                                                let signature =
1446                                                    Self::order_signature(&order_response.order);
1447                                                let duplicate_update = was_terminal
1448                                                    || context.last_order_signature
1449                                                        == Some(signature);
1450                                                let event = (!duplicate_update).then(|| {
1451                                                    parse_order_updated_with_client_order_id(
1452                                                        &order_response.order,
1453                                                        instrument,
1454                                                        account_id,
1455                                                        context.trader_id,
1456                                                        context.strategy_id,
1457                                                        context.client_order_id,
1458                                                        ts_init,
1459                                                    )
1460                                                });
1461                                                context.accepted = true;
1462                                                context.last_order_signature = Some(signature);
1463
1464                                                if was_terminal {
1465                                                    self.finish_order_context(
1466                                                        venue_order_id,
1467                                                        &context,
1468                                                    );
1469                                                } else {
1470                                                    self.bind_order_context(
1471                                                        venue_order_id,
1472                                                        context,
1473                                                    );
1474                                                }
1475                                                let outgoing = self.route_user_trades(
1476                                                    &order_response.trades,
1477                                                    ts_init,
1478                                                );
1479                                                self.pending_outgoing.extend(outgoing);
1480                                                return event.map(NautilusWsMessage::OrderUpdated);
1481                                            } else {
1482                                                log::warn!(
1483                                                    "Cannot create OrderUpdated: account_id not set"
1484                                                );
1485                                            }
1486                                        } else {
1487                                            log::warn!(
1488                                                "Instrument {instrument_name_ustr} not found in cache for edit response"
1489                                            );
1490                                        }
1491                                    }
1492                                    Err(e) => {
1493                                        log::error!(
1494                                            "Failed to parse edit response: request_id={request_id}, error={e}"
1495                                        );
1496                                    }
1497                                }
1498                            } else if let Some(error) = &response.error {
1499                                log::error!(
1500                                    "Order modify rejected: code={}, message={}, client_order_id={}",
1501                                    error.code,
1502                                    error.message,
1503                                    client_order_id
1504                                );
1505                                return Some(NautilusWsMessage::OrderModifyRejected(
1506                                    OrderModifyRejected::new(
1507                                        trader_id,
1508                                        strategy_id,
1509                                        instrument_id,
1510                                        client_order_id,
1511                                        ustr::ustr(&format!(
1512                                            "code={}: {}",
1513                                            error.code, error.message
1514                                        )),
1515                                        UUID4::new(),
1516                                        ts_init,
1517                                        ts_init,
1518                                        false,
1519                                        None, // venue_order_id not available
1520                                        self.account_id,
1521                                    ),
1522                                ));
1523                            }
1524                        }
1525                        PendingRequestType::GetOrderState {
1526                            client_order_id,
1527                            trader_id: _,
1528                            strategy_id: _,
1529                            instrument_id: _,
1530                        } => {
1531                            if let Some(result) = &response.result {
1532                                match serde_json::from_value::<DeribitOrderMsg>(result.clone()) {
1533                                    Ok(order_msg) => {
1534                                        log::debug!(
1535                                            "Order state received: venue_order_id={}, client_order_id={}, state={}",
1536                                            order_msg.order_id,
1537                                            client_order_id,
1538                                            order_msg.order_state
1539                                        );
1540
1541                                        // Convert to OrderStatusReport
1542                                        let instrument_name_ustr = order_msg.instrument_name;
1543
1544                                        if let Some(instrument) =
1545                                            self.instruments_cache.get(&instrument_name_ustr)
1546                                        {
1547                                            if let Some(account_id) = self.account_id {
1548                                                match parse_user_order_msg(
1549                                                    &order_msg, instrument, account_id, ts_init,
1550                                                ) {
1551                                                    Ok(report) => {
1552                                                        return Some(
1553                                                            NautilusWsMessage::OrderStatusReports(
1554                                                                vec![report],
1555                                                            ),
1556                                                        );
1557                                                    }
1558                                                    Err(e) => {
1559                                                        log::warn!(
1560                                                            "Failed to parse get_order_state response to report: {e}"
1561                                                        );
1562                                                    }
1563                                                }
1564                                            } else {
1565                                                log::warn!(
1566                                                    "Cannot create OrderStatusReport: account_id not set"
1567                                                );
1568                                            }
1569                                        } else {
1570                                            log::warn!(
1571                                                "Instrument {instrument_name_ustr} not found in cache for get_order_state response"
1572                                            );
1573                                        }
1574                                    }
1575                                    Err(e) => {
1576                                        log::error!(
1577                                            "Failed to parse get_order_state response: request_id={request_id}, error={e}"
1578                                        );
1579                                    }
1580                                }
1581                            } else if let Some(error) = &response.error {
1582                                log::error!(
1583                                    "Get order state failed: code={}, message={}, client_order_id={}",
1584                                    error.code,
1585                                    error.message,
1586                                    client_order_id
1587                                );
1588                            }
1589                        }
1590                    }
1591                } else if let Some(request_id) = response.id {
1592                    // Response with ID but no matching pending request
1593                    if let Some(error) = &response.error {
1594                        // Log orphaned error response with all available context
1595                        log::error!(
1596                            "Deribit error for unknown request: code={}, message={}, request_id={}, data={:?}",
1597                            error.code,
1598                            error.message,
1599                            request_id,
1600                            error.data
1601                        );
1602                        return Some(NautilusWsMessage::Error(DeribitWsError::DeribitError {
1603                            code: error.code,
1604                            message: error.message.clone(),
1605                        }));
1606                    } else {
1607                        // Success response but no pending request - likely already processed
1608                        log::debug!(
1609                            "Received response for unknown request_id={}, result present: {}",
1610                            request_id,
1611                            response.result.is_some()
1612                        );
1613                    }
1614                } else if let Some(error) = &response.error {
1615                    // Error response with no ID (shouldn't happen in JSON-RPC 2.0, but handle it)
1616                    log::error!(
1617                        "Deribit error with no request_id: code={}, message={}, data={:?}",
1618                        error.code,
1619                        error.message,
1620                        error.data
1621                    );
1622                    return Some(NautilusWsMessage::Error(DeribitWsError::DeribitError {
1623                        code: error.code,
1624                        message: error.message.clone(),
1625                    }));
1626                }
1627                None
1628            }
1629            DeribitWsMessage::Notification(notification) => {
1630                let channel = &notification.params.channel;
1631                let data = &notification.params.data;
1632
1633                // Determine channel type and parse accordingly
1634                if let Some(channel_type) = DeribitWsChannel::from_channel_string(channel) {
1635                    match channel_type {
1636                        DeribitWsChannel::Trades => {
1637                            // Parse trade messages
1638                            match serde_json::from_value::<Vec<DeribitTradeMsg>>(data.clone()) {
1639                                Ok(trades) => {
1640                                    log::debug!("Received {} trades", trades.len());
1641                                    let data_vec = parse_trades_data(
1642                                        &trades,
1643                                        &self.instruments_cache,
1644                                        ts_init,
1645                                    );
1646
1647                                    if data_vec.is_empty() && !trades.is_empty() {
1648                                        let missing: Vec<&Ustr> = trades
1649                                            .iter()
1650                                            .map(|t| &t.instrument_name)
1651                                            .filter(|name| {
1652                                                !self.instruments_cache.contains_key(name)
1653                                            })
1654                                            .collect();
1655
1656                                        if missing.is_empty() {
1657                                            log::warn!(
1658                                                "Received {} trades but parsed 0 (parse failures); cache size: {}",
1659                                                trades.len(),
1660                                                self.instruments_cache.len()
1661                                            );
1662                                        } else {
1663                                            log::warn!(
1664                                                "Trade message received but instrument(s) not found in cache: {:?} (cache size: {})",
1665                                                missing,
1666                                                self.instruments_cache.len()
1667                                            );
1668                                        }
1669                                    } else if !data_vec.is_empty() {
1670                                        log::debug!("Parsed {} trade ticks", data_vec.len());
1671                                        return Some(NautilusWsMessage::Data(data_vec));
1672                                    }
1673                                }
1674                                Err(e) => {
1675                                    log::warn!("Failed to deserialize trades: {e}");
1676                                }
1677                            }
1678                        }
1679                        DeribitWsChannel::Book => {
1680                            // Parse order book messages
1681                            match serde_json::from_value::<DeribitBookMsg>(data.clone()) {
1682                                Ok(book_msg) => {
1683                                    if let Some(instrument) =
1684                                        self.instruments_cache.get(&book_msg.instrument_name)
1685                                    {
1686                                        let inst_name = book_msg.instrument_name.to_string();
1687                                        let awaiting_resync =
1688                                            self.pending_book_resync.iter().any(|ch| {
1689                                                ch.starts_with("book.")
1690                                                    && ch
1691                                                        .split('.')
1692                                                        .nth(1)
1693                                                        .is_some_and(|s| s == inst_name)
1694                                            });
1695
1696                                        if awaiting_resync
1697                                            && book_msg.msg_type == DeribitBookMsgType::Change
1698                                        {
1699                                            // Drop deltas while awaiting resync snapshot
1700                                        } else if awaiting_resync
1701                                            && book_msg.msg_type == DeribitBookMsgType::Snapshot
1702                                        {
1703                                            self.pending_book_resync.retain(|ch| {
1704                                                !(ch.starts_with("book.")
1705                                                    && ch
1706                                                        .split('.')
1707                                                        .nth(1)
1708                                                        .is_some_and(|s| s == inst_name))
1709                                            });
1710                                            self.book_sequence.insert(
1711                                                book_msg.instrument_name,
1712                                                book_msg.change_id,
1713                                            );
1714
1715                                            match parse_book_msg(&book_msg, instrument, ts_init) {
1716                                                Ok(deltas) => {
1717                                                    return Some(NautilusWsMessage::Deltas(deltas));
1718                                                }
1719                                                Err(e) => {
1720                                                    log::warn!("Failed to parse book message: {e}");
1721                                                }
1722                                            }
1723                                        } else if book_msg.msg_type == DeribitBookMsgType::Change
1724                                            && let Some(prev_id) = book_msg.prev_change_id
1725                                            && let Some(&last_id) =
1726                                                self.book_sequence.get(&book_msg.instrument_name)
1727                                            && prev_id != last_id
1728                                        {
1729                                            log::error!(
1730                                                "Book sequence gap for {}: expected prev_change_id={}, was {} \
1731                                                - dropping delta, forcing resync",
1732                                                book_msg.instrument_name,
1733                                                last_id,
1734                                                prev_id
1735                                            );
1736                                            self.book_sequence.remove(&book_msg.instrument_name);
1737
1738                                            let book_channels: Vec<String> = self
1739                                                .subscriptions_state
1740                                                .all_topics()
1741                                                .into_iter()
1742                                                .filter(|t| {
1743                                                    t.starts_with("book.")
1744                                                        && t.split('.')
1745                                                            .nth(1)
1746                                                            .is_some_and(|s| s == inst_name)
1747                                                })
1748                                                .collect();
1749
1750                                            if !book_channels.is_empty() {
1751                                                for ch in &book_channels {
1752                                                    self.subscriptions_state.mark_failure(ch);
1753                                                }
1754                                                // Defer resubscribe until unsubscribe ack
1755                                                self.pending_book_resync
1756                                                    .extend(book_channels.clone());
1757                                                let _ =
1758                                                    self.handle_unsubscribe(book_channels).await;
1759                                            }
1760                                        } else {
1761                                            self.book_sequence.insert(
1762                                                book_msg.instrument_name,
1763                                                book_msg.change_id,
1764                                            );
1765
1766                                            match parse_book_msg(&book_msg, instrument, ts_init) {
1767                                                Ok(deltas) => {
1768                                                    return Some(NautilusWsMessage::Deltas(deltas));
1769                                                }
1770                                                Err(e) => {
1771                                                    log::warn!("Failed to parse book message: {e}");
1772                                                }
1773                                            }
1774                                        }
1775                                    } else {
1776                                        log::warn!(
1777                                            "Book message received but instrument '{}' not found in cache (cache size: {})",
1778                                            book_msg.instrument_name,
1779                                            self.instruments_cache.len()
1780                                        );
1781                                    }
1782                                }
1783                                Err(e) => {
1784                                    log::warn!(
1785                                        "Failed to deserialize book message: {e}, channel: {channel}"
1786                                    );
1787                                }
1788                            }
1789                        }
1790                        DeribitWsChannel::Ticker => {
1791                            match serde_json::from_value::<DeribitTickerMsg>(data.clone()) {
1792                                Ok(ticker_msg) => {
1793                                    if let Some(instrument) =
1794                                        self.instruments_cache.get(&ticker_msg.instrument_name)
1795                                    {
1796                                        // Emit OptionGreeks only if subscribed
1797                                        if self.option_greeks_subs.contains(&instrument.id())
1798                                            && let Some(option_greeks) =
1799                                                parse_ticker_to_option_greeks(
1800                                                    &ticker_msg,
1801                                                    instrument,
1802                                                    ts_init,
1803                                                )
1804                                        {
1805                                            let _ = self.out_tx.send(
1806                                                NautilusWsMessage::OptionGreeks(option_greeks),
1807                                            );
1808                                        }
1809
1810                                        let instrument_id = instrument.id();
1811                                        let mut data_vec = Vec::new();
1812
1813                                        // Emit MarkPriceUpdate only if subscribed
1814                                        if self.mark_price_subs.contains(&instrument_id) {
1815                                            match parse_ticker_to_mark_price(
1816                                                &ticker_msg,
1817                                                instrument,
1818                                                ts_init,
1819                                            ) {
1820                                                Ok(mark_price) => {
1821                                                    data_vec.push(Data::MarkPrice(mark_price));
1822                                                }
1823                                                Err(e) => {
1824                                                    log::warn!("Failed to parse mark price: {e}");
1825                                                }
1826                                            }
1827                                        }
1828
1829                                        // Emit IndexPriceUpdate only if subscribed
1830                                        if self.index_price_subs.contains(&instrument_id) {
1831                                            match parse_ticker_to_index_price(
1832                                                &ticker_msg,
1833                                                instrument,
1834                                                ts_init,
1835                                            ) {
1836                                                Ok(index_price) => {
1837                                                    data_vec.push(Data::IndexPrice(index_price));
1838                                                }
1839                                                Err(e) => {
1840                                                    log::warn!("Failed to parse index price: {e}");
1841                                                }
1842                                            }
1843                                        }
1844
1845                                        if !data_vec.is_empty() {
1846                                            return Some(NautilusWsMessage::Data(data_vec));
1847                                        }
1848                                    } else {
1849                                        log::warn!(
1850                                            "Ticker message received but instrument '{}' not found in cache (cache size: {})",
1851                                            ticker_msg.instrument_name,
1852                                            self.instruments_cache.len()
1853                                        );
1854                                    }
1855                                }
1856                                Err(e) => {
1857                                    log::warn!(
1858                                        "Failed to deserialize ticker message: {e}, channel: {channel}"
1859                                    );
1860                                }
1861                            }
1862                        }
1863                        DeribitWsChannel::Perpetual => {
1864                            // Parse perpetual channel for funding rate updates
1865                            // This channel is dedicated to perpetual instruments and provides
1866                            // the interest (funding) rate
1867                            match serde_json::from_value::<DeribitPerpetualMsg>(data.clone()) {
1868                                Ok(perpetual_msg) => {
1869                                    // Extract instrument name from channel: perpetual.{instrument}.{interval}
1870                                    let parts: Vec<&str> = channel.split('.').collect();
1871                                    if parts.len() >= 2 {
1872                                        let instrument_name = Ustr::from(parts[1]);
1873
1874                                        if let Some(instrument) =
1875                                            self.instruments_cache.get(&instrument_name)
1876                                        {
1877                                            let funding_rate = parse_perpetual_to_funding_rate(
1878                                                &perpetual_msg,
1879                                                instrument,
1880                                                ts_init,
1881                                            );
1882                                            return Some(NautilusWsMessage::FundingRates(vec![
1883                                                funding_rate,
1884                                            ]));
1885                                        } else {
1886                                            log::warn!(
1887                                                "Instrument {} not found in cache (cache size: {})",
1888                                                instrument_name,
1889                                                self.instruments_cache.len()
1890                                            );
1891                                        }
1892                                    }
1893                                }
1894                                Err(e) => {
1895                                    log::warn!(
1896                                        "Failed to deserialize perpetual message: {e}, data: {data}"
1897                                    );
1898                                }
1899                            }
1900                        }
1901                        DeribitWsChannel::Quote => {
1902                            // Parse quote messages
1903                            match serde_json::from_value::<DeribitQuoteMsg>(data.clone()) {
1904                                Ok(quote_msg) => {
1905                                    if let Some(instrument) =
1906                                        self.instruments_cache.get(&quote_msg.instrument_name)
1907                                    {
1908                                        match parse_quote_msg(&quote_msg, instrument, ts_init) {
1909                                            Ok(quote) => {
1910                                                return Some(NautilusWsMessage::Data(vec![
1911                                                    Data::Quote(quote),
1912                                                ]));
1913                                            }
1914                                            Err(e) => {
1915                                                log::warn!("Failed to parse quote message: {e}");
1916                                            }
1917                                        }
1918                                    } else {
1919                                        log::warn!(
1920                                            "Quote message received but instrument '{}' not found in cache (cache size: {})",
1921                                            quote_msg.instrument_name,
1922                                            self.instruments_cache.len()
1923                                        );
1924                                    }
1925                                }
1926                                Err(e) => {
1927                                    log::warn!(
1928                                        "Failed to deserialize quote message: {e}, channel: {channel}"
1929                                    );
1930                                }
1931                            }
1932                        }
1933                        DeribitWsChannel::VolatilityIndex => {
1934                            match serde_json::from_value::<DeribitVolatilityIndexMsg>(data.clone())
1935                            {
1936                                Ok(msg) => {
1937                                    let ts_event = UnixNanos::from(msg.timestamp * 1_000_000);
1938                                    let mut metadata = nautilus_core::Params::new();
1939                                    metadata.insert(
1940                                        "index_name".to_string(),
1941                                        serde_json::Value::String(msg.index_name.clone()),
1942                                    );
1943                                    let data_type = DataType::new(
1944                                        "DeribitVolatilityIndex",
1945                                        Some(metadata),
1946                                        None,
1947                                    );
1948
1949                                    let dvol = DeribitVolatilityIndex::new(
1950                                        msg.index_name,
1951                                        msg.volatility,
1952                                        ts_event,
1953                                        ts_init,
1954                                    );
1955
1956                                    return Some(NautilusWsMessage::Data(vec![Data::Custom(
1957                                        CustomData::new(Arc::new(dvol), data_type),
1958                                    )]));
1959                                }
1960                                Err(e) => {
1961                                    log::warn!("Failed to deserialize volatility index: {e}");
1962                                }
1963                            }
1964                        }
1965                        DeribitWsChannel::InstrumentState => {
1966                            match serde_json::from_value::<DeribitInstrumentStateMsg>(data.clone())
1967                            {
1968                                Ok(state_msg) => {
1969                                    log::debug!(
1970                                        "Instrument state change: {} -> {} (timestamp: {})",
1971                                        state_msg.instrument_name,
1972                                        state_msg.state,
1973                                        state_msg.timestamp
1974                                    );
1975
1976                                    let instrument_id = if let Some(instrument) =
1977                                        self.instruments_cache.get(&state_msg.instrument_name)
1978                                    {
1979                                        instrument.id()
1980                                    } else {
1981                                        log::debug!(
1982                                            "Instrument '{}' not in cache, constructing ID",
1983                                            state_msg.instrument_name
1984                                        );
1985                                        InstrumentId::new(
1986                                            Symbol::new(state_msg.instrument_name),
1987                                            *DERIBIT_VENUE,
1988                                        )
1989                                    };
1990
1991                                    let action = MarketStatusAction::from(state_msg.state);
1992                                    let is_trading =
1993                                        Some(state_msg.state == DeribitInstrumentState::Started);
1994                                    let ts_event = UnixNanos::from(state_msg.timestamp * 1_000_000);
1995                                    let status = InstrumentStatus::new(
1996                                        instrument_id,
1997                                        action,
1998                                        ts_event,
1999                                        ts_init,
2000                                        None,
2001                                        None,
2002                                        is_trading,
2003                                        None,
2004                                        None,
2005                                    );
2006                                    return Some(NautilusWsMessage::InstrumentStatus(status));
2007                                }
2008                                Err(e) => {
2009                                    log::warn!("Failed to parse instrument status message: {e}");
2010                                }
2011                            }
2012                        }
2013                        DeribitWsChannel::ChartTrades => {
2014                            // Parse chart.trades messages into Bar objects using emit-on-next pattern.
2015                            // Deribit sends updates for the current bar as it builds. We only emit
2016                            // a bar when we receive a bar with a different timestamp, confirming
2017                            // the previous bar is closed.
2018                            if let Ok(chart_msg) =
2019                                serde_json::from_value::<DeribitChartMsg>(data.clone())
2020                            {
2021                                // Extract instrument and resolution from channel
2022                                // Channel format: chart.trades.{instrument}.{resolution}
2023                                let parts: Vec<&str> = channel.split('.').collect();
2024                                if parts.len() >= 4 {
2025                                    let instrument_name = Ustr::from(parts[2]);
2026                                    let resolution = parts[3];
2027
2028                                    if let Some(instrument) =
2029                                        self.instruments_cache.get(&instrument_name)
2030                                    {
2031                                        let instrument_id = instrument.id();
2032
2033                                        match resolution_to_bar_type(instrument_id, resolution) {
2034                                            Ok(bar_type) => {
2035                                                let price_precision = instrument.price_precision();
2036                                                let size_precision = instrument.size_precision();
2037                                                let use_cost_for_volume =
2038                                                    use_cost_for_bar_volume(instrument);
2039
2040                                                match parse_chart_msg(
2041                                                    &chart_msg,
2042                                                    bar_type,
2043                                                    price_precision,
2044                                                    size_precision,
2045                                                    use_cost_for_volume,
2046                                                    self.bars_timestamp_on_close,
2047                                                    ts_init,
2048                                                ) {
2049                                                    Ok(new_bar) => {
2050                                                        // Check if we have a pending bar for this channel
2051                                                        let channel_key = channel.clone();
2052
2053                                                        if let Some(pending_bar) =
2054                                                            self.pending_bars.get(&channel_key)
2055                                                        {
2056                                                            // If new bar has different timestamp, the pending bar is closed
2057                                                            if new_bar.ts_event
2058                                                                != pending_bar.ts_event
2059                                                            {
2060                                                                let closed_bar = *pending_bar;
2061                                                                self.pending_bars
2062                                                                    .insert(channel_key, new_bar);
2063                                                                log::debug!(
2064                                                                    "Emitting closed bar: {closed_bar:?}"
2065                                                                );
2066                                                                return Some(
2067                                                                    NautilusWsMessage::Data(vec![
2068                                                                        Data::Bar(closed_bar),
2069                                                                    ]),
2070                                                                );
2071                                                            }
2072                                                            // Same timestamp - update pending bar with latest values
2073                                                            self.pending_bars
2074                                                                .insert(channel_key, new_bar);
2075                                                        } else {
2076                                                            // First bar for this channel - store as pending
2077                                                            self.pending_bars
2078                                                                .insert(channel_key, new_bar);
2079                                                        }
2080                                                    }
2081                                                    Err(e) => {
2082                                                        log::warn!(
2083                                                            "Failed to parse chart message to bar: {e}"
2084                                                        );
2085                                                    }
2086                                                }
2087                                            }
2088                                            Err(e) => {
2089                                                log::warn!(
2090                                                    "Failed to create BarType from resolution {resolution}: {e}"
2091                                                );
2092                                            }
2093                                        }
2094                                    } else {
2095                                        log::warn!(
2096                                            "Instrument {instrument_name} not found in cache for chart data"
2097                                        );
2098                                    }
2099                                }
2100                            }
2101                        }
2102                        DeribitWsChannel::UserOrders => {
2103                            // Handle both array and single object responses
2104                            let orders_result =
2105                                serde_json::from_value::<Vec<DeribitOrderMsg>>(data.clone())
2106                                    .or_else(|_| {
2107                                        serde_json::from_value::<DeribitOrderMsg>(data.clone())
2108                                            .map(|order| vec![order])
2109                                    });
2110
2111                            match orders_result {
2112                                Ok(orders) => {
2113                                    log::debug!("Received {} user order updates", orders.len());
2114
2115                                    // Require account_id for parsing
2116                                    let Some(account_id) = self.account_id else {
2117                                        log::warn!("Cannot parse user orders: account_id not set");
2118                                        return Some(NautilusWsMessage::Raw(data.clone()));
2119                                    };
2120
2121                                    let mut outgoing = Vec::new();
2122
2123                                    // Process each order and emit appropriate events
2124                                    for order in &orders {
2125                                        let venue_order_id_str = &order.order_id;
2126                                        let venue_order_id =
2127                                            VenueOrderId::new(venue_order_id_str.as_str());
2128                                        let instrument_name = order.instrument_name;
2129
2130                                        let Some(instrument) =
2131                                            self.instruments_cache.get(&instrument_name)
2132                                        else {
2133                                            log::warn!(
2134                                                "Instrument {instrument_name} not found in cache"
2135                                            );
2136                                            continue;
2137                                        };
2138
2139                                        let label_client_order_id = order
2140                                            .label
2141                                            .as_ref()
2142                                            .filter(|l| !l.is_empty())
2143                                            .map(ClientOrderId::new);
2144                                        let was_terminal = self
2145                                            .terminal_order_contexts
2146                                            .contains_key(&venue_order_id);
2147                                        let Some(mut context) = self.find_order_context(
2148                                            venue_order_id,
2149                                            label_client_order_id,
2150                                        ) else {
2151                                            match parse_user_order_msg(
2152                                                order, instrument, account_id, ts_init,
2153                                            ) {
2154                                                Ok(report) => outgoing.push(
2155                                                    NautilusWsMessage::OrderStatusReports(vec![
2156                                                        report,
2157                                                    ]),
2158                                                ),
2159                                                Err(e) => log::warn!(
2160                                                    "Failed to parse external order update: {e}"
2161                                                ),
2162                                            }
2163                                            continue;
2164                                        };
2165
2166                                        let signature = Self::order_signature(order);
2167
2168                                        // Determine event type based on order state
2169                                        let event_type = determine_order_event_type(
2170                                            &order.order_state,
2171                                            !context.accepted,
2172                                            order.replaced
2173                                                && context.last_order_signature != Some(signature),
2174                                        );
2175                                        context.last_order_signature = Some(signature);
2176
2177                                        let trader_id = context.trader_id;
2178                                        let strategy_id = context.strategy_id;
2179                                        let client_order_id = context.client_order_id;
2180
2181                                        match event_type {
2182                                            OrderEventType::Accepted => {
2183                                                // Skip if order already reached terminal state (race condition)
2184                                                if self
2185                                                    .terminal_order_contexts
2186                                                    .contains_key(&venue_order_id)
2187                                                {
2188                                                    log::debug!(
2189                                                        "Skipping OrderAccepted for terminal order: client_order_id={client_order_id}"
2190                                                    );
2191                                                    continue;
2192                                                }
2193
2194                                                let event =
2195                                                    parse_order_accepted_with_client_order_id(
2196                                                        order,
2197                                                        instrument,
2198                                                        account_id,
2199                                                        trader_id,
2200                                                        strategy_id,
2201                                                        client_order_id,
2202                                                        ts_init,
2203                                                    );
2204                                                context.accepted = true;
2205                                                self.bind_order_context(venue_order_id, context);
2206
2207                                                log::debug!(
2208                                                    "Emitting OrderAccepted: venue_order_id={venue_order_id}"
2209                                                );
2210                                                outgoing
2211                                                    .push(NautilusWsMessage::OrderAccepted(event));
2212                                            }
2213                                            OrderEventType::Canceled => {
2214                                                // Skip if already emitted from the cancel
2215                                                // response path
2216                                                if self
2217                                                    .terminal_order_contexts
2218                                                    .contains_key(&venue_order_id)
2219                                                {
2220                                                    log::trace!(
2221                                                        "Skipping duplicate OrderCanceled: client_order_id={client_order_id}"
2222                                                    );
2223                                                    continue;
2224                                                }
2225
2226                                                if !context.accepted {
2227                                                    outgoing
2228                                                        .push(NautilusWsMessage::OrderAccepted(
2229                                                        parse_order_accepted_with_client_order_id(
2230                                                            order,
2231                                                            instrument,
2232                                                            account_id,
2233                                                            trader_id,
2234                                                            strategy_id,
2235                                                            client_order_id,
2236                                                            ts_init,
2237                                                        ),
2238                                                    ));
2239                                                    context.accepted = true;
2240                                                }
2241
2242                                                let event =
2243                                                    parse_order_canceled_with_client_order_id(
2244                                                        order,
2245                                                        instrument,
2246                                                        account_id,
2247                                                        trader_id,
2248                                                        strategy_id,
2249                                                        client_order_id,
2250                                                        ts_init,
2251                                                    );
2252                                                log::debug!(
2253                                                    "Emitting OrderCanceled: venue_order_id={venue_order_id}"
2254                                                );
2255                                                self.finish_order_context(venue_order_id, &context);
2256                                                outgoing
2257                                                    .push(NautilusWsMessage::OrderCanceled(event));
2258                                            }
2259                                            OrderEventType::Expired => {
2260                                                if self
2261                                                    .terminal_order_contexts
2262                                                    .contains_key(&venue_order_id)
2263                                                {
2264                                                    log::trace!(
2265                                                        "Skipping duplicate OrderExpired: client_order_id={client_order_id}"
2266                                                    );
2267                                                    continue;
2268                                                }
2269
2270                                                if !context.accepted {
2271                                                    outgoing
2272                                                        .push(NautilusWsMessage::OrderAccepted(
2273                                                        parse_order_accepted_with_client_order_id(
2274                                                            order,
2275                                                            instrument,
2276                                                            account_id,
2277                                                            trader_id,
2278                                                            strategy_id,
2279                                                            client_order_id,
2280                                                            ts_init,
2281                                                        ),
2282                                                    ));
2283                                                    context.accepted = true;
2284                                                }
2285
2286                                                let event =
2287                                                    parse_order_expired_with_client_order_id(
2288                                                        order,
2289                                                        instrument,
2290                                                        account_id,
2291                                                        trader_id,
2292                                                        strategy_id,
2293                                                        client_order_id,
2294                                                        ts_init,
2295                                                    );
2296                                                log::debug!(
2297                                                    "Emitting OrderExpired: venue_order_id={venue_order_id}"
2298                                                );
2299                                                self.finish_order_context(venue_order_id, &context);
2300                                                outgoing
2301                                                    .push(NautilusWsMessage::OrderExpired(event));
2302                                            }
2303                                            OrderEventType::Updated => {
2304                                                if was_terminal {
2305                                                    log::trace!(
2306                                                        "Skipping amendment for terminal order: venue_order_id={venue_order_id}"
2307                                                    );
2308                                                    continue;
2309                                                }
2310
2311                                                if !context.accepted {
2312                                                    outgoing
2313                                                        .push(NautilusWsMessage::OrderAccepted(
2314                                                        parse_order_accepted_with_client_order_id(
2315                                                            order,
2316                                                            instrument,
2317                                                            account_id,
2318                                                            trader_id,
2319                                                            strategy_id,
2320                                                            client_order_id,
2321                                                            ts_init,
2322                                                        ),
2323                                                    ));
2324                                                    context.accepted = true;
2325                                                }
2326
2327                                                let event =
2328                                                    parse_order_updated_with_client_order_id(
2329                                                        order,
2330                                                        instrument,
2331                                                        account_id,
2332                                                        trader_id,
2333                                                        strategy_id,
2334                                                        client_order_id,
2335                                                        ts_init,
2336                                                    );
2337                                                self.bind_order_context(venue_order_id, context);
2338                                                log::debug!(
2339                                                    "Emitting OrderUpdated: venue_order_id={venue_order_id}"
2340                                                );
2341                                                outgoing
2342                                                    .push(NautilusWsMessage::OrderUpdated(event));
2343                                            }
2344                                            OrderEventType::None => {
2345                                                // Fills handled via user.trades, track terminal state
2346                                                // for race condition prevention
2347                                                if order.order_state == "filled" {
2348                                                    log::debug!(
2349                                                        "Recording terminal order: venue_order_id={venue_order_id}, state={}",
2350                                                        order.order_state
2351                                                    );
2352                                                    self.finish_order_context(
2353                                                        venue_order_id,
2354                                                        &context,
2355                                                    );
2356                                                } else if order.order_state == "rejected" {
2357                                                    log::debug!(
2358                                                        "Recording rejected order: venue_order_id={venue_order_id}"
2359                                                    );
2360                                                    self.finish_order_context(
2361                                                        venue_order_id,
2362                                                        &context,
2363                                                    );
2364                                                } else if was_terminal {
2365                                                    self.finish_order_context(
2366                                                        venue_order_id,
2367                                                        &context,
2368                                                    );
2369                                                } else {
2370                                                    log::trace!(
2371                                                        "No event to emit for order {}, state={}",
2372                                                        venue_order_id,
2373                                                        order.order_state
2374                                                    );
2375                                                    self.bind_order_context(
2376                                                        venue_order_id,
2377                                                        context,
2378                                                    );
2379                                                }
2380                                            }
2381                                        }
2382                                    }
2383
2384                                    if !outgoing.is_empty() {
2385                                        self.pending_outgoing.extend(outgoing);
2386                                    }
2387                                }
2388                                Err(e) => {
2389                                    log::warn!("Failed to deserialize user orders: {e}");
2390                                }
2391                            }
2392                        }
2393                        DeribitWsChannel::UserTrades => {
2394                            // Handle both array and single object responses
2395                            let trades_result =
2396                                serde_json::from_value::<Vec<DeribitUserTradeMsg>>(data.clone())
2397                                    .or_else(|_| {
2398                                        serde_json::from_value::<DeribitUserTradeMsg>(data.clone())
2399                                            .map(|trade| vec![trade])
2400                                    });
2401
2402                            match trades_result {
2403                                Ok(trades) => {
2404                                    log::debug!("Received {} user trade updates", trades.len());
2405                                    if self.account_id.is_none() {
2406                                        log::warn!("Cannot parse user trades: account_id not set");
2407                                        return Some(NautilusWsMessage::Raw(data.clone()));
2408                                    }
2409                                    let outgoing = self.route_user_trades(&trades, ts_init);
2410                                    if !outgoing.is_empty() {
2411                                        self.pending_outgoing.extend(outgoing);
2412                                    }
2413                                }
2414                                Err(e) => {
2415                                    log::warn!("Failed to deserialize user trades: {e}");
2416                                }
2417                            }
2418                        }
2419                        DeribitWsChannel::UserPortfolio => {
2420                            match serde_json::from_value::<DeribitPortfolioMsg>(data.clone()) {
2421                                Ok(portfolio) => {
2422                                    // Skip zero-balance currencies (common with cross-collateral)
2423                                    // Only check equity and balance - initial_margin can be non-zero
2424                                    // for all currencies when cross-collateral is enabled
2425                                    if portfolio.equity.is_zero() && portfolio.balance.is_zero() {
2426                                        log::trace!(
2427                                            "Skipping zero-balance portfolio for {}",
2428                                            portfolio.currency
2429                                        );
2430                                        return None;
2431                                    }
2432
2433                                    // Require account_id for parsing
2434                                    let Some(account_id) = self.account_id else {
2435                                        log::warn!("Cannot parse portfolio: account_id not set");
2436                                        return None;
2437                                    };
2438
2439                                    match parse_portfolio_to_account_state(
2440                                        &portfolio, account_id, ts_init,
2441                                    ) {
2442                                        Ok(account_state) => {
2443                                            // Check for duplicate per currency
2444                                            let currency_key = portfolio.currency.clone();
2445
2446                                            if let Some(last) =
2447                                                self.last_account_states.get(&currency_key)
2448                                                && account_state.has_same_balances_and_margins(last)
2449                                            {
2450                                                log::trace!(
2451                                                    "Skipping duplicate portfolio update for {}",
2452                                                    portfolio.currency
2453                                                );
2454                                                return None;
2455                                            }
2456
2457                                            self.last_account_states
2458                                                .insert(currency_key, account_state.clone());
2459                                            return Some(NautilusWsMessage::AccountState(
2460                                                account_state,
2461                                            ));
2462                                        }
2463                                        Err(e) => {
2464                                            log::warn!(
2465                                                "Failed to parse portfolio to AccountState: {e}"
2466                                            );
2467                                        }
2468                                    }
2469                                }
2470                                Err(e) => {
2471                                    log::warn!("Failed to deserialize portfolio: {e}");
2472                                }
2473                            }
2474                        }
2475                        _ => {
2476                            // Unhandled channel - return raw
2477                            log::trace!("Unhandled channel: {channel}");
2478                            return Some(NautilusWsMessage::Raw(data.clone()));
2479                        }
2480                    }
2481                } else {
2482                    log::trace!("Unknown channel: {channel}");
2483                    return Some(NautilusWsMessage::Raw(data.clone()));
2484                }
2485                None
2486            }
2487            DeribitWsMessage::Heartbeat(heartbeat) => {
2488                match heartbeat.heartbeat_type {
2489                    DeribitHeartbeatType::TestRequest => {
2490                        log::trace!(
2491                            "Received heartbeat test_request - responding with public/test"
2492                        );
2493
2494                        if let Err(e) = self.handle_heartbeat_test_request().await {
2495                            log::error!("Failed to respond to heartbeat test_request: {e}");
2496
2497                            // Return error to signal connection may be unhealthy
2498                            return Some(NautilusWsMessage::Error(DeribitWsError::Send(format!(
2499                                "Heartbeat response failed: {e}"
2500                            ))));
2501                        }
2502                    }
2503                    DeribitHeartbeatType::Heartbeat => {
2504                        log::trace!("Received heartbeat acknowledgment");
2505                    }
2506                }
2507                None
2508            }
2509            DeribitWsMessage::Error(err) => {
2510                log::error!("Deribit error {}: {}", err.code, err.message);
2511                Some(NautilusWsMessage::Error(DeribitWsError::DeribitError {
2512                    code: err.code,
2513                    message: err.message,
2514                }))
2515            }
2516            DeribitWsMessage::Reconnected => Some(NautilusWsMessage::Reconnected),
2517        }
2518    }
2519
2520    fn order_signature(order: &DeribitOrderMsg) -> OrderSignature {
2521        (order.amount, order.price, order.trigger_price)
2522    }
2523
2524    fn find_order_context(
2525        &self,
2526        venue_order_id: VenueOrderId,
2527        client_order_id: Option<ClientOrderId>,
2528    ) -> Option<OrderContext> {
2529        self.order_contexts
2530            .get(&venue_order_id)
2531            .or_else(|| self.terminal_order_contexts.get(&venue_order_id))
2532            .cloned()
2533            .or_else(|| {
2534                client_order_id.and_then(|client_order_id| {
2535                    self.submitted_order_contexts.get(&client_order_id).cloned()
2536                })
2537            })
2538    }
2539
2540    fn bind_order_context(&mut self, venue_order_id: VenueOrderId, context: OrderContext) {
2541        self.submitted_order_contexts
2542            .remove(&context.client_order_id);
2543        self.terminal_order_contexts.remove(&venue_order_id);
2544        self.order_contexts.insert(venue_order_id, context);
2545    }
2546
2547    fn finish_order_context(&mut self, venue_order_id: VenueOrderId, context: &OrderContext) {
2548        self.order_contexts.remove(&venue_order_id);
2549        self.submitted_order_contexts
2550            .remove(&context.client_order_id);
2551        self.terminal_order_contexts
2552            .insert(venue_order_id, context.clone());
2553    }
2554
2555    fn route_user_trades(
2556        &mut self,
2557        trades: &[DeribitUserTradeMsg],
2558        ts_init: UnixNanos,
2559    ) -> Vec<NautilusWsMessage> {
2560        let Some(account_id) = self.account_id else {
2561            log::warn!("Cannot parse user trades: account_id not set");
2562            return Vec::new();
2563        };
2564
2565        let mut outgoing = Vec::with_capacity(trades.len() + 1);
2566        let mut reports = Vec::new();
2567
2568        for trade in trades {
2569            let instrument_name = trade.instrument_name;
2570            let Some((report, quote_currency)) =
2571                self.instruments_cache
2572                    .get(&instrument_name)
2573                    .map(|instrument| {
2574                        (
2575                            parse_user_trade_msg(trade, instrument, account_id, ts_init),
2576                            instrument.quote_currency(),
2577                        )
2578                    })
2579            else {
2580                log::warn!("Instrument {instrument_name} not found in cache");
2581                continue;
2582            };
2583
2584            let report = match report {
2585                Ok(report) => report,
2586                Err(e) => {
2587                    log::warn!("Failed to parse trade {}: {e}", trade.trade_id);
2588                    continue;
2589                }
2590            };
2591            let venue_order_id = report.venue_order_id;
2592            let was_terminal = self.terminal_order_contexts.contains_key(&venue_order_id);
2593            let Some(mut context) = self.find_order_context(venue_order_id, report.client_order_id)
2594            else {
2595                log::debug!(
2596                    "Parsed external fill report: {} @ {}",
2597                    report.trade_id,
2598                    report.last_px
2599                );
2600                reports.push(report);
2601                continue;
2602            };
2603
2604            if !context.accepted {
2605                outgoing.push(NautilusWsMessage::OrderAccepted(OrderAccepted::new(
2606                    context.trader_id,
2607                    context.strategy_id,
2608                    context.instrument_id,
2609                    context.client_order_id,
2610                    venue_order_id,
2611                    account_id,
2612                    UUID4::new(),
2613                    report.ts_event,
2614                    report.ts_init,
2615                    false,
2616                )));
2617                context.accepted = true;
2618            }
2619
2620            log::debug!(
2621                "Parsed tracked fill event: {} @ {}",
2622                report.trade_id,
2623                report.last_px
2624            );
2625            outgoing.push(NautilusWsMessage::OrderFilled(OrderFilled::new(
2626                context.trader_id,
2627                context.strategy_id,
2628                context.instrument_id,
2629                context.client_order_id,
2630                venue_order_id,
2631                account_id,
2632                report.trade_id,
2633                context.order_side,
2634                context.order_type,
2635                report.last_qty,
2636                report.last_px,
2637                quote_currency,
2638                report.liquidity_side,
2639                UUID4::new(),
2640                report.ts_event,
2641                report.ts_init,
2642                false,
2643                report.venue_position_id,
2644                Some(report.commission),
2645                None,
2646            )));
2647
2648            if was_terminal || trade.state == "filled" {
2649                self.finish_order_context(venue_order_id, &context);
2650            } else {
2651                self.bind_order_context(venue_order_id, context);
2652            }
2653        }
2654
2655        if !reports.is_empty() {
2656            outgoing.push(NautilusWsMessage::FillReports(reports));
2657        }
2658        outgoing
2659    }
2660
2661    /// Main message processing loop.
2662    ///
2663    /// Returns `None` when the handler should stop.
2664    /// Messages that need client-side handling (e.g., Reconnected) are returned.
2665    /// Data messages are sent directly to `out_tx` for the user stream.
2666    pub async fn next(&mut self) -> Option<NautilusWsMessage> {
2667        loop {
2668            if let Some(msg) = self.pending_outgoing.pop_front() {
2669                match msg {
2670                    NautilusWsMessage::Reconnected
2671                    | NautilusWsMessage::Authenticated(_)
2672                    | NautilusWsMessage::AuthenticationFailed(_) => {
2673                        return Some(msg);
2674                    }
2675                    _ => {
2676                        let _ = self.out_tx.send(msg);
2677                        continue;
2678                    }
2679                }
2680            }
2681
2682            tokio::select! {
2683                // Process commands from client
2684                Some(cmd) = self.cmd_rx.recv() => {
2685                    self.process_command(cmd).await;
2686                }
2687                // Process raw WebSocket messages
2688                Some(msg) = self.raw_rx.recv() => {
2689                    match msg {
2690                        Message::Text(text) => {
2691                            if let Some(nautilus_msg) = self.process_raw_message(&text).await {
2692                                // Send data messages to user stream
2693                                match &nautilus_msg {
2694                                    NautilusWsMessage::Data(_)
2695                                    | NautilusWsMessage::Deltas(_)
2696                                    | NautilusWsMessage::Instrument(_)
2697                                    | NautilusWsMessage::InstrumentStatus(_)
2698                                    | NautilusWsMessage::OptionGreeks(_)
2699                                    | NautilusWsMessage::Raw(_)
2700                                    | NautilusWsMessage::Error(_) => {
2701                                        let _ = self.out_tx.send(nautilus_msg);
2702                                    }
2703                                    NautilusWsMessage::FundingRates(rates) => {
2704                                        let msg_to_send =
2705                                            NautilusWsMessage::FundingRates(rates.clone());
2706
2707                                        if let Err(e) = self.out_tx.send(msg_to_send) {
2708                                            log::error!("Failed to send funding rates: {e}");
2709                                        }
2710                                    }
2711                                    NautilusWsMessage::OrderStatusReports(_)
2712                                    | NautilusWsMessage::FillReports(_)
2713                                    | NautilusWsMessage::OrderFilled(_)
2714                                    | NautilusWsMessage::OrderAccepted(_)
2715                                    | NautilusWsMessage::OrderCanceled(_)
2716                                    | NautilusWsMessage::OrderExpired(_)
2717                                    | NautilusWsMessage::OrderUpdated(_)
2718                                    | NautilusWsMessage::OrderRejected(_)
2719                                    | NautilusWsMessage::OrderCancelRejected(_)
2720                                    | NautilusWsMessage::OrderModifyRejected(_)
2721                                    | NautilusWsMessage::AccountState(_) => {
2722                                        let _ = self.out_tx.send(nautilus_msg);
2723                                    }
2724                                    // Return messages that need client-side handling
2725                                    NautilusWsMessage::Reconnected
2726                                    | NautilusWsMessage::Authenticated(_)
2727                                    | NautilusWsMessage::AuthenticationFailed(_) => {
2728                                        return Some(nautilus_msg);
2729                                    }
2730                                }
2731                            }
2732                        }
2733                        Message::Ping(data) => {
2734                            // Respond to ping with pong
2735                            if let Some(client) = &self.inner {
2736                                let _ = client.send_pong(data.to_vec()).await;
2737                            }
2738                        }
2739                        Message::Close(_) => {
2740                            log::debug!("Received close frame");
2741                        }
2742                        _ => {}
2743                    }
2744                }
2745                // Check for stop signal
2746                () = tokio::time::sleep(tokio::time::Duration::from_millis(100)) => {
2747                    if self.signal.load(Ordering::Relaxed) {
2748                        log::debug!("Stop signal received");
2749                        return None;
2750                    }
2751                }
2752            }
2753        }
2754    }
2755}
2756
2757#[cfg(test)]
2758mod tests {
2759    use nautilus_model::{enums::LiquiditySide, instruments::Instrument, types::Money};
2760    use rstest::rstest;
2761
2762    use super::*;
2763    use crate::{
2764        common::{parse::parse_deribit_instrument_any, testing::load_test_json},
2765        http::models::{DeribitInstrument, DeribitJsonRpcResponse},
2766    };
2767
2768    fn routing_test_handler() -> DeribitWsFeedHandler {
2769        let signal = Arc::new(AtomicBool::new(false));
2770        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
2771        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
2772        let (out_tx, _out_rx) = tokio::sync::mpsc::unbounded_channel();
2773        let mut handler = DeribitWsFeedHandler::new(
2774            signal,
2775            cmd_rx,
2776            raw_rx,
2777            out_tx,
2778            AuthTracker::new(),
2779            SubscriptionState::new('.'),
2780            Arc::new(AtomicSet::new()),
2781            Arc::new(AtomicSet::new()),
2782            Arc::new(AtomicSet::new()),
2783            Some(AccountId::from("DERIBIT-001")),
2784            true,
2785            Arc::new(Mutex::new(Vec::new())),
2786        );
2787        let json = load_test_json("http_get_instruments.json");
2788        let response: DeribitJsonRpcResponse<Vec<DeribitInstrument>> =
2789            serde_json::from_str(&json).unwrap();
2790        let instrument = parse_deribit_instrument_any(
2791            &response.result.unwrap()[0],
2792            UnixNanos::default(),
2793            UnixNanos::default(),
2794        )
2795        .unwrap()
2796        .unwrap();
2797        handler
2798            .instruments_cache
2799            .insert(instrument.raw_symbol().inner(), instrument);
2800        handler
2801    }
2802
2803    fn order_data(order_state: &str, replaced: bool) -> serde_json::Value {
2804        serde_json::json!({
2805            "order_id": "ETH-584830574",
2806            "label": "O-19700101-000000-001-001-1",
2807            "instrument_name": "BTC-PERPETUAL",
2808            "direction": "buy",
2809            "order_type": "market",
2810            "order_state": order_state,
2811            "replaced": replaced,
2812            "price": 203.8,
2813            "amount": 2.0,
2814            "filled_amount": if order_state == "filled" { 2.0 } else { 0.0 },
2815            "average_price": 203.8,
2816            "creation_timestamp": 1_590_480_712_700_u64,
2817            "last_update_timestamp": 1_590_480_712_800_u64,
2818            "time_in_force": "good_til_cancelled",
2819            "commission": 0.00073602,
2820            "post_only": false,
2821            "reduce_only": false,
2822            "trigger_price": null,
2823            "trigger": null,
2824            "max_show": null,
2825            "api": true,
2826            "reject_reason": null,
2827            "cancel_reason": null
2828        })
2829    }
2830
2831    fn trade_data() -> serde_json::Value {
2832        serde_json::json!({
2833            "trade_id": "ETH-2696068",
2834            "order_id": "ETH-584830574",
2835            "instrument_name": "BTC-PERPETUAL",
2836            "direction": "buy",
2837            "price": 203.8,
2838            "amount": 2.0,
2839            "fee": 0.00073602,
2840            "fee_currency": "USDT",
2841            "timestamp": 1_590_480_712_800_u64,
2842            "trade_seq": 1_966_042_u64,
2843            "liquidity": "T",
2844            "order_type": "market",
2845            "index_price": 203.89,
2846            "mark_price": 203.78,
2847            "tick_direction": 3,
2848            "state": "filled",
2849            "label": "O-19700101-000000-001-001-1",
2850            "reduce_only": false,
2851            "post_only": false,
2852            "liquidation": null,
2853            "profit_loss": null
2854        })
2855    }
2856
2857    fn subscription(channel: &str, data: &serde_json::Value) -> String {
2858        serde_json::json!({
2859            "jsonrpc": "2.0",
2860            "method": "subscription",
2861            "params": { "channel": channel, "data": data }
2862        })
2863        .to_string()
2864    }
2865
2866    #[rstest]
2867    #[tokio::test]
2868    async fn tracked_fast_fill_synthesizes_accepted_before_filled() {
2869        let mut handler = routing_test_handler();
2870        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
2871        let instrument_id = InstrumentId::from("BTC-PERPETUAL.DERIBIT");
2872        handler.submitted_order_contexts.insert(
2873            client_order_id,
2874            OrderContext {
2875                client_order_id,
2876                trader_id: TraderId::from("TRADER-001"),
2877                strategy_id: StrategyId::from("S-001"),
2878                instrument_id,
2879                order_side: OrderSide::Buy,
2880                order_type: OrderType::Market,
2881                accepted: false,
2882                last_order_signature: None,
2883            },
2884        );
2885        handler.pending_requests.insert(
2886            42,
2887            PendingRequestType::Buy {
2888                client_order_id,
2889                trader_id: TraderId::from("TRADER-001"),
2890                strategy_id: StrategyId::from("S-001"),
2891                instrument_id,
2892                order_side: OrderSide::Buy,
2893                order_type: OrderType::Market,
2894            },
2895        );
2896        let response = serde_json::json!({
2897            "jsonrpc": "2.0",
2898            "id": 42,
2899            "result": { "order": order_data("filled", false), "trades": [] }
2900        });
2901
2902        handler.process_raw_message(&response.to_string()).await;
2903        assert!(handler.pending_outgoing.is_empty());
2904        assert!(
2905            !handler
2906                .terminal_order_contexts
2907                .get(&VenueOrderId::from("ETH-584830574"))
2908                .unwrap()
2909                .accepted
2910        );
2911
2912        handler
2913            .process_raw_message(&subscription(
2914                "user.trades.any.any.raw",
2915                &serde_json::json!([trade_data()]),
2916            ))
2917            .await;
2918
2919        assert!(matches!(
2920            handler.pending_outgoing.pop_front().unwrap(),
2921            NautilusWsMessage::OrderAccepted(event)
2922                if event.client_order_id == client_order_id
2923        ));
2924        assert!(matches!(
2925            handler.pending_outgoing.pop_front().unwrap(),
2926            NautilusWsMessage::OrderFilled(event)
2927                if event.client_order_id == client_order_id
2928                    && event.trade_id.to_string() == "ETH-2696068"
2929                    && event.order_side == OrderSide::Buy
2930                    && event.order_type == OrderType::Market
2931                    && event.last_qty.to_string() == "2"
2932                    && event.last_px.to_string() == "203.8"
2933                    && event.liquidity_side == LiquiditySide::Taker
2934                    && event.commission == Some(Money::from("0.00073602 USDT"))
2935        ));
2936        assert!(handler.pending_outgoing.is_empty());
2937
2938        handler
2939            .process_raw_message(&subscription(
2940                "user.trades.any.any.raw",
2941                &serde_json::json!([trade_data()]),
2942            ))
2943            .await;
2944
2945        assert!(matches!(
2946            handler.pending_outgoing.pop_front().unwrap(),
2947            NautilusWsMessage::OrderFilled(event)
2948                if event.client_order_id == client_order_id
2949        ));
2950        assert!(handler.pending_outgoing.is_empty());
2951    }
2952
2953    #[rstest]
2954    #[tokio::test]
2955    async fn submit_response_trades_use_tracked_event_path() {
2956        let mut handler = routing_test_handler();
2957        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
2958        let context = OrderContext {
2959            client_order_id,
2960            trader_id: TraderId::from("TRADER-001"),
2961            strategy_id: StrategyId::from("S-001"),
2962            instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
2963            order_side: OrderSide::Buy,
2964            order_type: OrderType::Market,
2965            accepted: false,
2966            last_order_signature: None,
2967        };
2968        handler
2969            .submitted_order_contexts
2970            .insert(client_order_id, context.clone());
2971        handler.pending_requests.insert(
2972            42,
2973            PendingRequestType::Buy {
2974                client_order_id,
2975                trader_id: context.trader_id,
2976                strategy_id: context.strategy_id,
2977                instrument_id: context.instrument_id,
2978                order_side: context.order_side,
2979                order_type: context.order_type,
2980            },
2981        );
2982        let response = serde_json::json!({
2983            "jsonrpc": "2.0",
2984            "id": 42,
2985            "result": { "order": order_data("filled", false), "trades": [trade_data()] }
2986        });
2987
2988        handler.process_raw_message(&response.to_string()).await;
2989
2990        assert!(matches!(
2991            handler.pending_outgoing.pop_front().unwrap(),
2992            NautilusWsMessage::OrderAccepted(event)
2993                if event.client_order_id == client_order_id
2994        ));
2995        assert!(matches!(
2996            handler.pending_outgoing.pop_front().unwrap(),
2997            NautilusWsMessage::OrderFilled(event)
2998                if event.client_order_id == client_order_id
2999                    && event.trade_id.to_string() == "ETH-2696068"
3000        ));
3001        assert!(handler.pending_outgoing.is_empty());
3002
3003        handler
3004            .process_raw_message(&subscription(
3005                "user.trades.any.any.raw",
3006                &serde_json::json!([trade_data()]),
3007            ))
3008            .await;
3009
3010        assert!(matches!(
3011            handler.pending_outgoing.pop_front().unwrap(),
3012            NautilusWsMessage::OrderFilled(event)
3013                if event.client_order_id == client_order_id
3014        ));
3015        assert!(handler.pending_outgoing.is_empty());
3016    }
3017
3018    #[rstest]
3019    #[tokio::test]
3020    async fn reconnect_preserves_submit_identity_before_response() {
3021        let mut handler = routing_test_handler();
3022        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3023        handler.submitted_order_contexts.insert(
3024            client_order_id,
3025            OrderContext {
3026                client_order_id,
3027                trader_id: TraderId::from("TRADER-001"),
3028                strategy_id: StrategyId::from("S-001"),
3029                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3030                order_side: OrderSide::Buy,
3031                order_type: OrderType::Market,
3032                accepted: false,
3033                last_order_signature: None,
3034            },
3035        );
3036        handler
3037            .pending_requests
3038            .insert(42, PendingRequestType::Test);
3039
3040        handler.clear_state();
3041        handler
3042            .process_raw_message(&subscription(
3043                "user.orders.any.any.raw",
3044                &serde_json::json!([order_data("open", false)]),
3045            ))
3046            .await;
3047
3048        assert!(handler.pending_requests.is_empty());
3049        assert!(matches!(
3050            handler.pending_outgoing.pop_front().unwrap(),
3051            NautilusWsMessage::OrderAccepted(event)
3052                if event.client_order_id == client_order_id
3053        ));
3054        assert!(
3055            handler
3056                .order_contexts
3057                .get(&VenueOrderId::from("ETH-584830574"))
3058                .unwrap()
3059                .accepted
3060        );
3061    }
3062
3063    #[rstest]
3064    #[case("cancelled")]
3065    #[case("expired")]
3066    #[tokio::test]
3067    async fn tracked_terminal_order_synthesizes_accepted_first(#[case] order_state: &str) {
3068        let mut handler = routing_test_handler();
3069        let venue_order_id = VenueOrderId::from("ETH-584830574");
3070        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3071        handler.order_contexts.insert(
3072            venue_order_id,
3073            OrderContext {
3074                client_order_id,
3075                trader_id: TraderId::from("TRADER-001"),
3076                strategy_id: StrategyId::from("S-001"),
3077                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3078                order_side: OrderSide::Buy,
3079                order_type: OrderType::Market,
3080                accepted: false,
3081                last_order_signature: None,
3082            },
3083        );
3084
3085        handler
3086            .process_raw_message(&subscription(
3087                "user.orders.any.any.raw",
3088                &serde_json::json!([order_data(order_state, false)]),
3089            ))
3090            .await;
3091
3092        let accepted = handler.pending_outgoing.pop_front().unwrap();
3093        let terminal = handler.pending_outgoing.pop_front().unwrap();
3094        assert!(matches!(
3095            accepted,
3096            NautilusWsMessage::OrderAccepted(event)
3097                if event.client_order_id == client_order_id
3098        ));
3099        assert!(
3100            matches!(
3101                (order_state, terminal),
3102                ("cancelled", NautilusWsMessage::OrderCanceled(_))
3103                    | ("expired", NautilusWsMessage::OrderExpired(_))
3104            ),
3105            "unexpected terminal message for {order_state}",
3106        );
3107        assert!(handler.pending_outgoing.is_empty());
3108        assert!(!handler.order_contexts.contains_key(&venue_order_id));
3109    }
3110
3111    #[rstest]
3112    #[case("filled")]
3113    #[case("open")]
3114    #[tokio::test]
3115    async fn late_fill_after_cancel_stays_on_tracked_event_path(#[case] trade_state: &str) {
3116        let mut handler = routing_test_handler();
3117        let venue_order_id = VenueOrderId::from("ETH-584830574");
3118        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3119        handler.order_contexts.insert(
3120            venue_order_id,
3121            OrderContext {
3122                client_order_id,
3123                trader_id: TraderId::from("TRADER-001"),
3124                strategy_id: StrategyId::from("S-001"),
3125                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3126                order_side: OrderSide::Buy,
3127                order_type: OrderType::Market,
3128                accepted: true,
3129                last_order_signature: None,
3130            },
3131        );
3132
3133        handler
3134            .process_raw_message(&subscription(
3135                "user.orders.any.any.raw",
3136                &serde_json::json!([order_data("cancelled", false)]),
3137            ))
3138            .await;
3139        let mut trade = trade_data();
3140        trade["state"] = serde_json::Value::String(trade_state.to_string());
3141        handler
3142            .process_raw_message(&subscription(
3143                "user.trades.any.any.raw",
3144                &serde_json::json!([trade]),
3145            ))
3146            .await;
3147
3148        assert!(matches!(
3149            handler.pending_outgoing.pop_front().unwrap(),
3150            NautilusWsMessage::OrderCanceled(event)
3151                if event.client_order_id == client_order_id
3152        ));
3153        assert!(matches!(
3154            handler.pending_outgoing.pop_front().unwrap(),
3155            NautilusWsMessage::OrderFilled(event)
3156                if event.client_order_id == client_order_id
3157        ));
3158        assert!(handler.pending_outgoing.is_empty());
3159        assert!(!handler.order_contexts.contains_key(&venue_order_id));
3160        assert!(
3161            handler
3162                .terminal_order_contexts
3163                .contains_key(&venue_order_id)
3164        );
3165    }
3166
3167    #[rstest]
3168    #[tokio::test]
3169    async fn tracked_order_uses_stored_client_id_when_label_changes() {
3170        let mut handler = routing_test_handler();
3171        let venue_order_id = VenueOrderId::from("ETH-584830574");
3172        let client_order_id = ClientOrderId::from("ORIGINAL-CLIENT-ID");
3173        handler.order_contexts.insert(
3174            venue_order_id,
3175            OrderContext {
3176                client_order_id,
3177                trader_id: TraderId::from("TRADER-001"),
3178                strategy_id: StrategyId::from("S-001"),
3179                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3180                order_side: OrderSide::Buy,
3181                order_type: OrderType::Market,
3182                accepted: false,
3183                last_order_signature: None,
3184            },
3185        );
3186        let mut order = order_data("open", false);
3187        order["label"] = serde_json::Value::Null;
3188
3189        handler
3190            .process_raw_message(&subscription(
3191                "user.orders.any.any.raw",
3192                &serde_json::json!([order]),
3193            ))
3194            .await;
3195
3196        assert!(matches!(
3197            handler.pending_outgoing.pop_front().unwrap(),
3198            NautilusWsMessage::OrderAccepted(event)
3199                if event.client_order_id == client_order_id
3200        ));
3201        assert!(handler.pending_outgoing.is_empty());
3202    }
3203
3204    #[rstest]
3205    #[tokio::test]
3206    async fn subscription_accept_before_submit_response_is_not_duplicated() {
3207        let mut handler = routing_test_handler();
3208        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3209        handler.submitted_order_contexts.insert(
3210            client_order_id,
3211            OrderContext {
3212                client_order_id,
3213                trader_id: TraderId::from("TRADER-001"),
3214                strategy_id: StrategyId::from("S-001"),
3215                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3216                order_side: OrderSide::Buy,
3217                order_type: OrderType::Market,
3218                accepted: false,
3219                last_order_signature: None,
3220            },
3221        );
3222        handler.pending_requests.insert(
3223            42,
3224            PendingRequestType::Buy {
3225                client_order_id,
3226                trader_id: TraderId::from("TRADER-001"),
3227                strategy_id: StrategyId::from("S-001"),
3228                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3229                order_side: OrderSide::Buy,
3230                order_type: OrderType::Market,
3231            },
3232        );
3233
3234        handler
3235            .process_raw_message(&subscription(
3236                "user.orders.any.any.raw",
3237                &serde_json::json!([order_data("open", false)]),
3238            ))
3239            .await;
3240        let accepted = handler.pending_outgoing.pop_front().unwrap();
3241        let response = serde_json::json!({
3242            "jsonrpc": "2.0",
3243            "id": 42,
3244            "result": { "order": order_data("open", false), "trades": [] }
3245        });
3246        let duplicate = handler.process_raw_message(&response.to_string()).await;
3247        handler.clear_state();
3248        handler
3249            .process_raw_message(&subscription(
3250                "user.orders.any.any.raw",
3251                &serde_json::json!([order_data("open", false)]),
3252            ))
3253            .await;
3254
3255        assert!(matches!(accepted, NautilusWsMessage::OrderAccepted(_)));
3256        assert!(duplicate.is_none());
3257        assert!(handler.pending_outgoing.is_empty());
3258        assert!(
3259            handler
3260                .order_contexts
3261                .get(&VenueOrderId::from("ETH-584830574"))
3262                .unwrap()
3263                .accepted
3264        );
3265    }
3266
3267    #[rstest]
3268    #[tokio::test]
3269    async fn untracked_live_order_and_fill_use_report_paths() {
3270        let mut handler = routing_test_handler();
3271
3272        handler
3273            .process_raw_message(&subscription(
3274                "user.orders.any.any.raw",
3275                &serde_json::json!([order_data("open", false)]),
3276            ))
3277            .await;
3278        handler
3279            .process_raw_message(&subscription(
3280                "user.trades.any.any.raw",
3281                &serde_json::json!([trade_data()]),
3282            ))
3283            .await;
3284
3285        assert!(matches!(
3286            handler.pending_outgoing.pop_front().unwrap(),
3287            NautilusWsMessage::OrderStatusReports(reports) if reports.len() == 1
3288        ));
3289        assert!(matches!(
3290            handler.pending_outgoing.pop_front().unwrap(),
3291            NautilusWsMessage::FillReports(reports) if reports.len() == 1
3292        ));
3293        assert!(handler.pending_outgoing.is_empty());
3294    }
3295
3296    #[rstest]
3297    #[tokio::test]
3298    async fn edit_response_without_tracked_context_uses_report_path() {
3299        let mut handler = routing_test_handler();
3300        handler.pending_requests.insert(
3301            42,
3302            PendingRequestType::Edit {
3303                client_order_id: ClientOrderId::from("UNKNOWN-CLIENT-ID"),
3304                trader_id: TraderId::from("TRADER-001"),
3305                strategy_id: StrategyId::from("S-001"),
3306                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3307            },
3308        );
3309        let response = serde_json::json!({
3310            "jsonrpc": "2.0",
3311            "id": 42,
3312            "result": { "order": order_data("open", true), "trades": [trade_data()] }
3313        });
3314
3315        let message = handler.process_raw_message(&response.to_string()).await;
3316        let fill = handler.pending_outgoing.pop_front().unwrap();
3317
3318        assert!(matches!(
3319            message,
3320            Some(NautilusWsMessage::OrderStatusReports(reports)) if reports.len() == 1
3321        ));
3322        assert!(matches!(
3323            fill,
3324            NautilusWsMessage::FillReports(reports) if reports.len() == 1
3325        ));
3326        assert!(handler.order_contexts.is_empty());
3327        assert!(handler.pending_outgoing.is_empty());
3328    }
3329
3330    #[rstest]
3331    #[tokio::test]
3332    async fn tracked_edit_response_routes_fill_and_deduplicates_subscription_echo() {
3333        let mut handler = routing_test_handler();
3334        let venue_order_id = VenueOrderId::from("ETH-584830574");
3335        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3336        handler.order_contexts.insert(
3337            venue_order_id,
3338            OrderContext {
3339                client_order_id,
3340                trader_id: TraderId::from("TRADER-001"),
3341                strategy_id: StrategyId::from("S-001"),
3342                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3343                order_side: OrderSide::Buy,
3344                order_type: OrderType::Market,
3345                accepted: true,
3346                last_order_signature: None,
3347            },
3348        );
3349        handler.pending_requests.insert(
3350            42,
3351            PendingRequestType::Edit {
3352                client_order_id,
3353                trader_id: TraderId::from("TRADER-001"),
3354                strategy_id: StrategyId::from("S-001"),
3355                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3356            },
3357        );
3358        let response = serde_json::json!({
3359            "jsonrpc": "2.0",
3360            "id": 42,
3361            "result": { "order": order_data("open", true), "trades": [trade_data()] }
3362        });
3363
3364        let updated = handler.process_raw_message(&response.to_string()).await;
3365        let filled = handler.pending_outgoing.pop_front().unwrap();
3366        handler
3367            .process_raw_message(&subscription(
3368                "user.orders.any.any.raw",
3369                &serde_json::json!([order_data("open", true)]),
3370            ))
3371            .await;
3372
3373        assert!(matches!(
3374            updated,
3375            Some(NautilusWsMessage::OrderUpdated(event))
3376                if event.client_order_id == client_order_id
3377        ));
3378        assert!(matches!(
3379            filled,
3380            NautilusWsMessage::OrderFilled(event)
3381                if event.client_order_id == client_order_id
3382        ));
3383        assert!(handler.pending_outgoing.is_empty());
3384        assert!(!handler.order_contexts.contains_key(&venue_order_id));
3385        assert!(
3386            handler
3387                .terminal_order_contexts
3388                .get(&venue_order_id)
3389                .unwrap()
3390                .last_order_signature
3391                .is_some(),
3392        );
3393    }
3394
3395    #[rstest]
3396    #[tokio::test]
3397    async fn subscription_edit_echo_before_response_is_not_duplicated() {
3398        let mut handler = routing_test_handler();
3399        let venue_order_id = VenueOrderId::from("ETH-584830574");
3400        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3401        handler.order_contexts.insert(
3402            venue_order_id,
3403            OrderContext {
3404                client_order_id,
3405                trader_id: TraderId::from("TRADER-001"),
3406                strategy_id: StrategyId::from("S-001"),
3407                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3408                order_side: OrderSide::Buy,
3409                order_type: OrderType::Market,
3410                accepted: true,
3411                last_order_signature: None,
3412            },
3413        );
3414        handler.pending_requests.insert(
3415            42,
3416            PendingRequestType::Edit {
3417                client_order_id,
3418                trader_id: TraderId::from("TRADER-001"),
3419                strategy_id: StrategyId::from("S-001"),
3420                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3421            },
3422        );
3423        handler
3424            .process_raw_message(&subscription(
3425                "user.orders.any.any.raw",
3426                &serde_json::json!([order_data("open", true)]),
3427            ))
3428            .await;
3429        let updated = handler.pending_outgoing.pop_front().unwrap();
3430        let response = serde_json::json!({
3431            "jsonrpc": "2.0",
3432            "id": 42,
3433            "result": { "order": order_data("open", true), "trades": [trade_data()] }
3434        });
3435
3436        let duplicate = handler.process_raw_message(&response.to_string()).await;
3437        let filled = handler.pending_outgoing.pop_front().unwrap();
3438
3439        assert!(matches!(
3440            updated,
3441            NautilusWsMessage::OrderUpdated(event)
3442                if event.client_order_id == client_order_id
3443        ));
3444        assert!(duplicate.is_none());
3445        assert!(matches!(
3446            filled,
3447            NautilusWsMessage::OrderFilled(event)
3448                if event.client_order_id == client_order_id
3449        ));
3450        assert!(handler.pending_outgoing.is_empty());
3451    }
3452
3453    #[rstest]
3454    #[tokio::test]
3455    async fn delayed_edit_response_keeps_partial_fill_terminal() {
3456        let mut handler = routing_test_handler();
3457        let venue_order_id = VenueOrderId::from("ETH-584830574");
3458        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3459        handler.order_contexts.insert(
3460            venue_order_id,
3461            OrderContext {
3462                client_order_id,
3463                trader_id: TraderId::from("TRADER-001"),
3464                strategy_id: StrategyId::from("S-001"),
3465                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3466                order_side: OrderSide::Buy,
3467                order_type: OrderType::Market,
3468                accepted: true,
3469                last_order_signature: None,
3470            },
3471        );
3472        handler
3473            .process_raw_message(&subscription(
3474                "user.trades.any.any.raw",
3475                &serde_json::json!([trade_data()]),
3476            ))
3477            .await;
3478        handler.pending_outgoing.clear();
3479        handler.pending_requests.insert(
3480            42,
3481            PendingRequestType::Edit {
3482                client_order_id,
3483                trader_id: TraderId::from("TRADER-001"),
3484                strategy_id: StrategyId::from("S-001"),
3485                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3486            },
3487        );
3488        let mut partial_trade = trade_data();
3489        partial_trade["state"] = serde_json::Value::String("open".to_string());
3490        let response = serde_json::json!({
3491            "jsonrpc": "2.0",
3492            "id": 42,
3493            "result": { "order": order_data("open", true), "trades": [partial_trade] }
3494        });
3495
3496        let updated = handler.process_raw_message(&response.to_string()).await;
3497        let filled = handler.pending_outgoing.pop_front().unwrap();
3498
3499        assert!(updated.is_none());
3500        assert!(matches!(
3501            filled,
3502            NautilusWsMessage::OrderFilled(event)
3503                if event.client_order_id == client_order_id
3504        ));
3505        assert!(!handler.order_contexts.contains_key(&venue_order_id));
3506        assert!(
3507            handler
3508                .terminal_order_contexts
3509                .contains_key(&venue_order_id)
3510        );
3511        assert!(handler.pending_outgoing.is_empty());
3512    }
3513
3514    #[rstest]
3515    #[tokio::test]
3516    async fn tracked_replaced_order_emits_updated_event() {
3517        let mut handler = routing_test_handler();
3518        let venue_order_id = VenueOrderId::from("ETH-584830574");
3519        let client_order_id = ClientOrderId::from("O-19700101-000000-001-001-1");
3520        handler.order_contexts.insert(
3521            venue_order_id,
3522            OrderContext {
3523                client_order_id,
3524                trader_id: TraderId::from("TRADER-001"),
3525                strategy_id: StrategyId::from("S-001"),
3526                instrument_id: InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
3527                order_side: OrderSide::Buy,
3528                order_type: OrderType::Market,
3529                accepted: true,
3530                last_order_signature: None,
3531            },
3532        );
3533
3534        handler
3535            .process_raw_message(&subscription(
3536                "user.orders.any.any.raw",
3537                &serde_json::json!([order_data("open", true)]),
3538            ))
3539            .await;
3540
3541        assert!(matches!(
3542            handler.pending_outgoing.pop_front().unwrap(),
3543            NautilusWsMessage::OrderUpdated(event)
3544                if event.client_order_id == client_order_id
3545        ));
3546
3547        let mut partial_fill = order_data("open", true);
3548        partial_fill["filled_amount"] = serde_json::json!(1.0);
3549        partial_fill["last_update_timestamp"] = serde_json::json!(1_590_480_712_900_u64);
3550        handler
3551            .process_raw_message(&subscription(
3552                "user.orders.any.any.raw",
3553                &serde_json::json!([partial_fill]),
3554            ))
3555            .await;
3556
3557        assert!(handler.pending_outgoing.is_empty());
3558    }
3559}