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