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