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