1use std::{future::Future, time::Duration};
19
20use anyhow::Context;
21use async_trait::async_trait;
22use futures_util::{StreamExt, pin_mut};
23use nautilus_common::{
24 clients::ExecutionClient,
25 live::runner::get_exec_event_sender,
26 messages::execution::{
27 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
28 GenerateFillReportsBuilder, GenerateOrderStatusReport, GenerateOrderStatusReports,
29 GenerateOrderStatusReportsBuilder, GeneratePositionStatusReports,
30 GeneratePositionStatusReportsBuilder, ModifyOrder, QueryAccount, QueryOrder, SubmitOrder,
31 SubmitOrderList,
32 },
33};
34use nautilus_core::{
35 DurationNanos, Params, UnixNanos,
36 time::{AtomicTime, get_atomic_clock_realtime},
37};
38use nautilus_live::{
39 ExecutionClientCore, ExecutionEventEmitter, SocketControl,
40 task::{TaskGroup, TaskGroupGuard},
41};
42use nautilus_model::{
43 accounts::AccountAny,
44 enums::{AccountType, OmsType, OrderType, TimeInForce},
45 events::OrderEventAny,
46 identifiers::{AccountId, ClientId, Venue},
47 orders::{Order, OrderAny},
48 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
49 types::{AccountBalance, MarginBalance},
50};
51
52use crate::{
53 common::{
54 consts::{DERIBIT_VENUE, DERIBIT_WS_HEARTBEAT_SECS},
55 enums::resolve_trigger_type,
56 },
57 config::DeribitExecutionClientConfig,
58 http::{client::DeribitHttpClient, models::DeribitCurrency, query::GetOrderStateParams},
59 websocket::{
60 auth::DERIBIT_EXECUTION_SESSION_NAME,
61 client::DeribitWebSocketClient,
62 messages::{DeribitOrderParams, NautilusWsMessage},
63 parse::parse_user_order_msg,
64 },
65};
66
67#[derive(Debug)]
69pub struct DeribitExecutionClient {
70 core: ExecutionClientCore,
71 clock: &'static AtomicTime,
72 config: DeribitExecutionClientConfig,
73 emitter: ExecutionEventEmitter,
74 http_client: DeribitHttpClient,
75 ws_client: DeribitWebSocketClient,
76 session_tasks: TaskGroup,
77 pending_tasks: TaskGroup,
78}
79
80impl DeribitExecutionClient {
81 pub fn new(
87 core: ExecutionClientCore,
88 config: DeribitExecutionClientConfig,
89 ) -> anyhow::Result<Self> {
90 let api_key = config
91 .api_key
92 .as_ref()
93 .map(|value| value.expose_secret().to_owned());
94 let api_secret = config
95 .api_secret
96 .as_ref()
97 .map(|value| value.expose_secret().to_owned());
98 let proxy_url = config
99 .proxy_url
100 .as_ref()
101 .map(|value| value.expose_secret().to_owned());
102 let http_client = if config.has_api_credentials() {
103 DeribitHttpClient::new_with_env(
104 api_key.clone(),
105 api_secret.clone(),
106 config.base_url_http.clone(),
107 config.environment,
108 config.http_timeout_secs,
109 config.max_retries,
110 config.retry_delay_initial_ms,
111 config.retry_delay_max_ms,
112 proxy_url.clone(),
113 )?
114 } else {
115 DeribitHttpClient::new(
116 config.base_url_http.clone(),
117 config.environment,
118 config.http_timeout_secs,
119 config.max_retries,
120 config.retry_delay_initial_ms,
121 config.retry_delay_max_ms,
122 proxy_url.clone(),
123 )?
124 };
125
126 let mut ws_client = DeribitWebSocketClient::new(
127 config.base_url_ws.clone(),
128 api_key,
129 api_secret,
130 DERIBIT_WS_HEARTBEAT_SECS,
131 config.auth_timeout_secs,
132 config.environment,
133 config.transport_backend,
134 proxy_url,
135 )
136 .context("failed to create WebSocket client for execution")?
137 .with_socket_control(SocketControl::new(
138 core.client_id,
139 Some(*DERIBIT_VENUE),
140 "deribit-user-streams",
141 ));
142 ws_client.set_account_id(core.account_id);
144
145 let clock = get_atomic_clock_realtime();
146 let emitter = ExecutionEventEmitter::new(
147 clock,
148 core.trader_id,
149 core.account_id,
150 AccountType::Margin,
151 None,
152 );
153
154 let session_tasks = TaskGroup::new();
155 let pending_tasks = TaskGroup::new();
156
157 Ok(Self {
158 core,
159 clock,
160 config,
161 emitter,
162 http_client,
163 ws_client,
164 session_tasks,
165 pending_tasks,
166 })
167 }
168
169 fn spawn_task<F>(&self, description: &'static str, fut: F)
171 where
172 F: Future<Output = anyhow::Result<()>> + Send + 'static,
173 {
174 let future = async move {
175 if let Err(e) = fut.await {
176 log::warn!("{description} failed: {e:?}");
177 }
178 };
179
180 if let Err(e) = self.pending_tasks.spawn(future) {
181 log::warn!("Skipping Deribit {description} after shutdown began: {e}");
182 }
183 }
184
185 fn abort_pending_tasks(&self) {
187 self.pending_tasks.begin_shutdown();
188 }
189
190 fn abort_session_tasks(&self) {
191 self.session_tasks.begin_shutdown();
192 self.ws_client.begin_shutdown();
193 }
194
195 async fn await_pending_tasks(&self) -> anyhow::Result<()> {
196 self.pending_tasks.begin_shutdown();
197 self.pending_tasks
198 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
199 .await
200 .map_err(|e| anyhow::anyhow!("Failed to terminate Deribit execution tasks: {e}"))?;
201 Ok(())
202 }
203
204 async fn await_session_tasks(&self) -> anyhow::Result<()> {
205 self.session_tasks.begin_shutdown();
206 self.session_tasks
207 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
208 .await
209 .map_err(|e| anyhow::anyhow!("Failed to terminate Deribit session tasks: {e}"))?;
210 Ok(())
211 }
212
213 async fn teardown_partial_connect(&self) -> anyhow::Result<()> {
214 self.abort_session_tasks();
215 self.abort_pending_tasks();
216
217 let mut errors = Vec::new();
218 if let Err(e) = self.ws_client.close().await {
219 errors.push(format!("WebSocket shutdown failed: {e}"));
220 }
221 let (session_result, pending_result) =
222 tokio::join!(self.await_session_tasks(), self.await_pending_tasks());
223
224 if let Err(e) = session_result {
225 errors.push(e.to_string());
226 }
227
228 if let Err(e) = pending_result {
229 errors.push(e.to_string());
230 }
231 self.core.set_disconnected();
232
233 if errors.is_empty() {
234 Ok(())
235 } else {
236 anyhow::bail!(errors.join("; "))
237 }
238 }
239
240 fn build_order_params(order: &dyn Order) -> anyhow::Result<DeribitOrderParams> {
242 let order_type = match order.order_type() {
243 OrderType::Limit => "limit",
244 OrderType::Market => "market",
245 OrderType::StopLimit => "stop_limit",
246 OrderType::StopMarket => "stop_market",
247 OrderType::LimitIfTouched => "take_limit",
248 OrderType::MarketIfTouched => "take_market",
249 other => {
250 anyhow::bail!("Unsupported order type {other:?} for Deribit");
251 }
252 }
253 .to_string();
254
255 let time_in_force = if matches!(
256 order.order_type(),
257 OrderType::Market | OrderType::StopMarket | OrderType::MarketIfTouched
258 ) {
259 None
261 } else {
262 Some(
263 match order.time_in_force() {
264 TimeInForce::Gtc => "good_til_cancelled",
265 TimeInForce::Ioc => "immediate_or_cancel",
266 TimeInForce::Fok => "fill_or_kill",
267 TimeInForce::Gtd => {
268 if order.expire_time().is_some() {
269 log::warn!(
270 "Deribit GTD orders expire at 8:00 UTC only - custom expire_time is ignored. \
271 For custom expiry times, use managed GTD with emulation_trigger"
272 );
273 }
274 "good_til_day"
275 }
276 other => {
277 anyhow::bail!("Unsupported time_in_force {other:?} for Deribit");
278 }
279 }
280 .to_string(),
281 )
282 };
283
284 let valid_until = None;
287
288 let trigger = resolve_trigger_type(order.trigger_type());
289
290 Ok(DeribitOrderParams {
291 instrument_name: order.instrument_id().symbol.to_string(),
292 amount: order.quantity().as_decimal(),
293 order_type,
294 label: Some(order.client_order_id().to_string()),
295 price: order.price().map(|p| p.as_decimal()),
296 time_in_force,
297 post_only: if order.is_post_only() {
298 Some(true)
299 } else {
300 None
301 },
302 reject_post_only: if order.is_post_only() {
303 Some(true)
304 } else {
305 None
306 },
307 reduce_only: if order.is_reduce_only() {
308 Some(true)
309 } else {
310 None
311 },
312 trigger_price: order.trigger_price().map(|p| p.as_decimal()),
313 trigger,
314 max_show: None,
315 valid_until,
316 })
317 }
318
319 fn submit_single_order(&self, order: &OrderAny, task_name: &'static str) {
323 if order.is_closed() {
324 log::warn!("Cannot submit closed order {}", order.client_order_id());
325 return;
326 }
327
328 let params = match Self::build_order_params(order) {
329 Ok(params) => params,
330 Err(e) => {
331 let ts_event = self.clock.get_time_ns();
332 self.emitter.emit_order_rejected_event(
333 order.strategy_id(),
334 order.instrument_id(),
335 order.client_order_id(),
336 &format!("{e}"),
337 ts_event,
338 false,
339 );
340 return;
341 }
342 };
343 let client_order_id = order.client_order_id();
344 let trader_id = order.trader_id();
345 let strategy_id = order.strategy_id();
346 let instrument_id = order.instrument_id();
347 let order_side = order.order_side();
348
349 log::debug!("OrderSubmitted client_order_id={client_order_id}");
350 self.emitter.emit_order_submitted(order);
351
352 let ws_client = self.ws_client.clone();
353
354 self.spawn_task(task_name, async move {
355 let result = ws_client
356 .submit_order(
357 order_side,
358 params,
359 client_order_id,
360 trader_id,
361 strategy_id,
362 instrument_id,
363 )
364 .await;
365
366 if let Err(e) = result {
367 log::error!(
368 "Submit order request failed: task={task_name}, client_order_id={client_order_id}, error={e}"
369 );
370 return Err(e.into());
371 }
372
373 Ok(())
374 });
375 }
376
377 fn spawn_stream_handler(
379 &self,
380 stream: impl futures_util::Stream<Item = NautilusWsMessage> + Send + 'static,
381 ) -> anyhow::Result<()> {
382 let emitter = self.emitter.clone();
383
384 self.session_tasks.spawn(async move {
385 pin_mut!(stream);
386 while let Some(message) = stream.next().await {
387 dispatch_ws_message(message, &emitter);
388 }
389 })?;
390
391 log::debug!("WebSocket stream handler started");
392 Ok(())
393 }
394}
395
396#[async_trait(?Send)]
397impl ExecutionClient for DeribitExecutionClient {
398 fn is_connected(&self) -> bool {
399 self.core.is_connected()
400 }
401
402 fn client_id(&self) -> ClientId {
403 self.core.client_id
404 }
405
406 fn account_id(&self) -> AccountId {
407 self.core.account_id
408 }
409
410 fn venue(&self) -> Venue {
411 *DERIBIT_VENUE
412 }
413
414 fn oms_type(&self) -> OmsType {
415 self.core.oms_type
416 }
417
418 fn get_account(&self) -> Option<AccountAny> {
419 self.core.cache().account_owned(&self.core.account_id)
420 }
421
422 fn generate_account_state(
423 &self,
424 balances: Vec<AccountBalance>,
425 margins: Vec<MarginBalance>,
426 reported: bool,
427 ts_event: UnixNanos,
428 info: Option<Params>,
429 ) -> anyhow::Result<()> {
430 self.emitter
431 .emit_account_state(balances, margins, reported, ts_event, info);
432 Ok(())
433 }
434
435 fn start(&mut self) -> anyhow::Result<()> {
436 if self.core.is_started() {
437 return Ok(());
438 }
439
440 let sender = get_exec_event_sender();
441 self.emitter.set_sender(sender);
442 self.core.set_started();
443
444 log::info!(
445 "Started: client_id={}, account_id={}, account_type={:?}, product_types={:?}, environment={}",
446 self.core.client_id,
447 self.core.account_id,
448 self.core.account_type,
449 self.config.product_types,
450 self.config.environment
451 );
452 Ok(())
453 }
454
455 fn stop(&mut self) -> anyhow::Result<()> {
456 if self.core.is_stopped() {
457 return Ok(());
458 }
459
460 self.core.set_stopped();
461 self.core.set_disconnected();
462 self.abort_session_tasks();
463 self.abort_pending_tasks();
464 log::info!("Stopped: client_id={}", self.core.client_id);
465 Ok(())
466 }
467
468 async fn connect(&mut self) -> anyhow::Result<()> {
469 if self.core.is_connected() && self.pending_tasks.is_open() && self.session_tasks.is_open()
470 {
471 return Ok(());
472 }
473
474 if !self.pending_tasks.is_open() {
475 self.await_pending_tasks().await?;
476 self.pending_tasks
477 .start_generation()
478 .map_err(|e| anyhow::anyhow!("Failed to start Deribit task generation: {e}"))?;
479 }
480
481 if !self.session_tasks.is_open() || !self.session_tasks.is_empty() {
482 self.abort_session_tasks();
483
484 if self.ws_client.is_active() {
485 self.ws_client
486 .close()
487 .await
488 .context("failed to close stale Deribit WebSocket")?;
489 }
490 self.await_session_tasks().await?;
491 self.session_tasks
492 .start_generation()
493 .map_err(|e| anyhow::anyhow!("Failed to start Deribit session generation: {e}"))?;
494 } else if self.ws_client.is_active() {
495 self.ws_client
496 .close()
497 .await
498 .context("failed to close stale Deribit WebSocket")?;
499 }
500 let ws_client = self.ws_client.clone();
501 let setup_guard =
502 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
503 ws_client.begin_shutdown();
504 });
505
506 if !self.config.has_api_credentials() {
508 anyhow::bail!("Missing API credentials; set Deribit environment variables");
509 }
510
511 self.ws_client.set_account_id(self.core.account_id);
513
514 if !self.core.instruments_initialized() {
516 for product_type in &self.config.product_types {
517 let instruments = self
518 .http_client
519 .request_instruments(DeribitCurrency::ANY, Some(*product_type))
520 .await
521 .with_context(|| {
522 format!("failed to request instruments for {product_type:?}")
523 })?;
524
525 if instruments.is_empty() {
526 log::warn!("No instruments returned for {product_type:?}");
527 continue;
528 }
529
530 log::debug!("Fetched {} {product_type:?} instruments", instruments.len());
531 self.ws_client.cache_instruments(&instruments);
532 self.http_client.cache_instruments(&instruments);
533 }
534 self.core.set_instruments_initialized();
535 }
536
537 let account_state = self
539 .http_client
540 .request_account_state(self.core.account_id)
541 .await
542 .context("failed to request account state")?;
543
544 self.emitter.send_account_state(account_state);
545
546 let session_result = async {
547 self.ws_client
548 .connect()
549 .await
550 .context("failed to connect WebSocket client for execution")?;
551
552 self.ws_client
553 .authenticate_session(DERIBIT_EXECUTION_SESSION_NAME)
554 .await
555 .map_err(|e| anyhow::anyhow!("failed to authenticate WebSocket session: {e}"))?;
556
557 log::debug!("WebSocket client authenticated for execution");
558
559 self.ws_client
561 .subscribe_user_orders()
562 .await
563 .map_err(|e| anyhow::anyhow!("failed to subscribe to user orders: {e}"))?;
564 self.ws_client
565 .subscribe_user_trades()
566 .await
567 .map_err(|e| anyhow::anyhow!("failed to subscribe to user trades: {e}"))?;
568 self.ws_client
569 .subscribe_user_portfolio()
570 .await
571 .map_err(|e| anyhow::anyhow!("failed to subscribe to user portfolio: {e}"))?;
572
573 if let Err(e) = self.ws_client.wait_for_subscriptions_confirmed(30.0).await {
574 let _ = self.ws_client.unsubscribe_user_orders().await;
576 let _ = self.ws_client.unsubscribe_user_trades().await;
577 let _ = self.ws_client.unsubscribe_user_portfolio().await;
578 anyhow::bail!("subscription confirmation failed: {e}");
579 }
580
581 log::debug!("Subscribed to user order, trade, and portfolio updates");
582
583 let stream = self.ws_client.stream()?;
585 self.spawn_stream_handler(stream)?;
586
587 Ok::<(), anyhow::Error>(())
588 }
589 .await;
590
591 if let Err(e) = session_result {
592 if let Err(teardown_error) = self.teardown_partial_connect().await {
593 return Err(e.context(format!(
594 "Deribit execution startup teardown failed: {teardown_error}"
595 )));
596 }
597 return Err(e);
598 }
599
600 self.core.set_connected();
601 setup_guard.disarm();
602 log::info!("Connected: client_id={}", self.core.client_id);
603 Ok(())
604 }
605
606 async fn disconnect(&mut self) -> anyhow::Result<()> {
607 self.teardown_partial_connect().await?;
608 log::info!("Disconnected: client_id={}", self.core.client_id);
609 Ok(())
610 }
611
612 async fn generate_order_status_report(
613 &self,
614 cmd: &GenerateOrderStatusReport,
615 ) -> anyhow::Result<Option<OrderStatusReport>> {
616 if let Some(venue_order_id) = &cmd.venue_order_id {
618 let params = GetOrderStateParams {
619 order_id: venue_order_id.to_string(),
620 };
621 let ts_init = self.clock.get_time_ns();
622
623 match self.http_client.inner.get_order_state(params).await {
624 Ok(response) => {
625 if let Some(order) = response.result {
626 let symbol = order.instrument_name;
627 if let Some(instrument) = self.http_client.get_instrument(&symbol) {
628 let report = parse_user_order_msg(
629 &order,
630 &instrument,
631 self.core.account_id,
632 ts_init,
633 )?;
634 return Ok(Some(report));
635 } else {
636 log::warn!(
637 "Instrument {} not in cache for order {}",
638 order.instrument_name,
639 order.order_id
640 );
641 }
642 }
643 }
644 Err(e) => {
645 log::warn!("Failed to get order state: {e}");
646 }
647 }
648 return Ok(None);
649 }
650
651 if let Some(client_order_id) = &cmd.client_order_id {
653 let reports = self
654 .http_client
655 .request_order_status_reports(
656 self.core.account_id,
657 cmd.instrument_id,
658 None,
659 None,
660 false, )
662 .await?;
663
664 for report in reports {
666 if report.client_order_id == Some(*client_order_id) {
667 return Ok(Some(report));
668 }
669 }
670 }
671
672 Ok(None)
673 }
674
675 async fn generate_order_status_reports(
676 &self,
677 cmd: &GenerateOrderStatusReports,
678 ) -> anyhow::Result<Vec<OrderStatusReport>> {
679 self.http_client
680 .request_order_status_reports(
681 self.core.account_id,
682 cmd.instrument_id,
683 cmd.start,
684 cmd.end,
685 cmd.open_only,
686 )
687 .await
688 }
689
690 async fn generate_fill_reports(
691 &self,
692 cmd: GenerateFillReports,
693 ) -> anyhow::Result<Vec<FillReport>> {
694 let mut reports = self
695 .http_client
696 .request_fill_reports(self.core.account_id, cmd.instrument_id, cmd.start, cmd.end)
697 .await?;
698
699 if let Some(venue_order_id) = &cmd.venue_order_id {
701 reports.retain(|r| r.venue_order_id == *venue_order_id);
702 }
703
704 Ok(reports)
705 }
706
707 async fn generate_position_status_reports(
708 &self,
709 cmd: &GeneratePositionStatusReports,
710 ) -> anyhow::Result<Vec<PositionStatusReport>> {
711 self.http_client
712 .request_position_status_reports(self.core.account_id, cmd.instrument_id)
713 .await
714 }
715
716 async fn generate_mass_status(
717 &self,
718 lookback_mins: Option<u64>,
719 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
720 log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
721 let ts_now = self.clock.get_time_ns();
722 let start = lookback_mins
723 .map(DurationNanos::try_from_mins)
724 .transpose()?
725 .map(|lookback| ts_now.saturating_sub(lookback));
726
727 let order_cmd = GenerateOrderStatusReportsBuilder::default()
728 .ts_init(ts_now)
729 .open_only(false) .start(start)
731 .build()
732 .context("Failed to build GenerateOrderStatusReports")?;
733
734 let fill_cmd = GenerateFillReportsBuilder::default()
735 .ts_init(ts_now)
736 .start(start)
737 .build()
738 .context("Failed to build GenerateFillReports")?;
739
740 let position_cmd = GeneratePositionStatusReportsBuilder::default()
741 .ts_init(ts_now)
742 .start(start)
743 .build()
744 .context("Failed to build GeneratePositionStatusReports")?;
745
746 let (order_reports, fill_reports, position_reports) = tokio::try_join!(
747 self.generate_order_status_reports(&order_cmd),
748 self.generate_fill_reports(fill_cmd),
749 self.generate_position_status_reports(&position_cmd),
750 )?;
751
752 log::info!("Received {} OrderStatusReports", order_reports.len());
753 log::info!("Received {} FillReports", fill_reports.len());
754 log::info!("Received {} PositionReports", position_reports.len());
755
756 let mut mass_status = ExecutionMassStatus::new(
757 self.core.client_id,
758 self.core.account_id,
759 *DERIBIT_VENUE,
760 ts_now,
761 None,
762 );
763
764 mass_status.add_order_reports(order_reports);
765 mass_status.add_fill_reports(fill_reports);
766 mass_status.add_position_reports(position_reports);
767
768 Ok(Some(mass_status))
769 }
770
771 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
772 let http_client = self.http_client.clone();
773 let account_id = self.core.account_id;
774 let emitter = self.emitter.clone();
775
776 self.spawn_task("query_account", async move {
777 let account_state = http_client
778 .request_account_state(account_id)
779 .await
780 .context("failed to query account state (check API credentials are valid)")?;
781
782 emitter.send_account_state(account_state);
783 Ok(())
784 });
785
786 Ok(())
787 }
788
789 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
790 let ws_client = self.ws_client.clone();
791
792 let order_id = cmd
794 .venue_order_id
795 .as_ref()
796 .ok_or_else(|| anyhow::anyhow!("venue_order_id required for query_order"))?
797 .to_string();
798
799 let client_order_id = cmd.client_order_id;
800 let trader_id = cmd.trader_id;
801 let strategy_id = cmd.strategy_id;
802 let instrument_id = cmd.instrument_id;
803
804 log::debug!("Querying order state: order_id={order_id}, client_order_id={client_order_id}");
805
806 self.spawn_task("query_order", async move {
809 ws_client
810 .query_order(
811 &order_id,
812 client_order_id,
813 trader_id,
814 strategy_id,
815 instrument_id,
816 )
817 .await
818 .map_err(|e| anyhow::anyhow!("Query order state failed: {e}"))?;
819 Ok(())
820 });
821
822 Ok(())
823 }
824
825 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
826 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
827 self.submit_single_order(&order, "submit_order");
828 Ok(())
829 }
830
831 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
832 if cmd.order_list.client_order_ids.is_empty() {
833 log::debug!("submit_order_list called with empty order list");
834 return Ok(());
835 }
836
837 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
838
839 log::debug!(
840 "Submitting order list {} with {} orders for instrument={}",
841 cmd.order_list.id,
842 orders.len(),
843 cmd.instrument_id
844 );
845
846 for order in &orders {
849 self.submit_single_order(order, "submit_order_list_item");
850 }
851
852 Ok(())
853 }
854
855 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
856 let ws_client = self.ws_client.clone();
857
858 let order_id = match cmd.venue_order_id.as_ref() {
860 Some(venue_order_id) => venue_order_id.to_string(),
861 None => {
862 return reject_modify_command(
863 &self.emitter,
864 self.clock,
865 &cmd,
866 "venue_order_id required for modify_order",
867 );
868 }
869 };
870
871 let quantity = if let Some(qty) = cmd.quantity {
873 qty
874 } else {
875 let cache = self.core.cache();
877 match cache.order(&cmd.client_order_id) {
878 Some(order) => order.quantity(),
879 None => {
880 return reject_modify_command(
881 &self.emitter,
882 self.clock,
883 &cmd,
884 &format!("Order not found: {}", cmd.client_order_id),
885 );
886 }
887 }
888 };
889
890 let price = match cmd.price {
891 Some(price) => price,
892 None => {
893 return reject_modify_command(
894 &self.emitter,
895 self.clock,
896 &cmd,
897 "price required for modify_order",
898 );
899 }
900 };
901
902 let client_order_id = cmd.client_order_id;
903 let trader_id = cmd.trader_id;
904 let strategy_id = cmd.strategy_id;
905 let instrument_id = cmd.instrument_id;
906
907 log::debug!(
908 "Modifying order: order_id={order_id}, quantity={quantity}, price={price}, client_order_id={client_order_id}"
909 );
910
911 self.spawn_task("modify_order", async move {
913 if let Err(e) = ws_client
914 .modify_order(
915 &order_id,
916 quantity,
917 price,
918 client_order_id,
919 trader_id,
920 strategy_id,
921 instrument_id,
922 )
923 .await
924 {
925 log::error!(
926 "Modify order failed: order_id={order_id}, client_order_id={client_order_id}, error={e}"
927 );
928 anyhow::bail!("Modify order failed: {e}");
929 }
930 Ok(())
931 });
932
933 Ok(())
934 }
935
936 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
937 let ws_client = self.ws_client.clone();
938
939 let order_id = match cmd.venue_order_id.as_ref() {
941 Some(venue_order_id) => venue_order_id.to_string(),
942 None => {
943 log::warn!(
944 "Cannot cancel order {} - no venue_order_id",
945 cmd.client_order_id
946 );
947 return Ok(());
948 }
949 };
950
951 let client_order_id = cmd.client_order_id;
952 let trader_id = cmd.trader_id;
953 let strategy_id = cmd.strategy_id;
954 let instrument_id = cmd.instrument_id;
955
956 log::debug!("Canceling order: order_id={order_id}, client_order_id={client_order_id}");
957
958 self.spawn_task("cancel_order", async move {
960 if let Err(e) = ws_client
961 .cancel_order(
962 &order_id,
963 client_order_id,
964 trader_id,
965 strategy_id,
966 instrument_id,
967 )
968 .await
969 {
970 log::error!(
971 "Cancel order failed: order_id={order_id}, client_order_id={client_order_id}, error={e}"
972 );
973 anyhow::bail!("Cancel order failed: {e}");
974 }
975 Ok(())
976 });
977
978 Ok(())
979 }
980
981 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
982 let instrument_id = cmd.instrument_id;
983
984 let Some(order_side) = cmd.order_side else {
986 log::debug!(
987 "Cancelling all orders: instrument={instrument_id}, order_side=None (bulk)"
988 );
989
990 let ws_client = self.ws_client.clone();
991 self.spawn_task("cancel_all_orders", async move {
992 if let Err(e) = ws_client.cancel_all_orders(instrument_id, None).await {
993 log::error!("Cancel all orders failed for instrument {instrument_id}: {e}");
994 anyhow::bail!("Cancel all orders failed: {e}");
995 }
996 Ok(())
997 });
998
999 return Ok(());
1000 };
1001
1002 log::debug!(
1005 "Cancelling orders by side: instrument={instrument_id}, order_side={order_side}"
1006 );
1007
1008 let orders_to_cancel: Vec<_> = {
1009 let cache = self.core.cache();
1010 let open_orders = cache.orders_open(None, Some(&instrument_id), None, None, None);
1011
1012 open_orders
1013 .into_iter()
1014 .filter(|order| order.order_side() == order_side)
1015 .filter_map(|order| {
1016 let venue_order_id = order.venue_order_id()?;
1017 Some((
1018 venue_order_id.to_string(),
1019 order.client_order_id(),
1020 order.instrument_id(),
1021 Some(venue_order_id),
1022 ))
1023 })
1024 .collect()
1025 };
1026
1027 if orders_to_cancel.is_empty() {
1028 log::debug!("No open {order_side} orders to cancel for {instrument_id}");
1029 return Ok(());
1030 }
1031
1032 log::debug!(
1033 "Cancelling {} {order_side} orders for {instrument_id}",
1034 orders_to_cancel.len(),
1035 );
1036
1037 for (venue_order_id_str, client_order_id, order_instrument_id, _venue_order_id) in
1039 orders_to_cancel
1040 {
1041 let ws_client = self.ws_client.clone();
1042 let trader_id = cmd.trader_id;
1043 let strategy_id = cmd.strategy_id;
1044
1045 self.spawn_task("cancel_order_by_side", async move {
1046 if let Err(e) = ws_client
1047 .cancel_order(
1048 &venue_order_id_str,
1049 client_order_id,
1050 trader_id,
1051 strategy_id,
1052 order_instrument_id,
1053 )
1054 .await
1055 {
1056 log::error!(
1057 "Cancel order failed: order_id={venue_order_id_str}, client_order_id={client_order_id}, error={e}"
1058 );
1059 }
1060 Ok(())
1061 });
1062 }
1063
1064 Ok(())
1065 }
1066
1067 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1068 if cmd.cancels.is_empty() {
1069 log::debug!("batch_cancel_orders called with empty cancels list");
1070 return Ok(());
1071 }
1072
1073 log::debug!(
1074 "Batch cancelling {} orders for instrument={}",
1075 cmd.cancels.len(),
1076 cmd.instrument_id
1077 );
1078
1079 for cancel in &cmd.cancels {
1082 let order_id = match &cancel.venue_order_id {
1083 Some(id) => id.to_string(),
1084 None => {
1085 log::warn!(
1086 "Cannot cancel order {} - no venue_order_id",
1087 cancel.client_order_id
1088 );
1089 continue;
1090 }
1091 };
1092
1093 let ws_client = self.ws_client.clone();
1094 let client_order_id = cancel.client_order_id;
1095 let trader_id = cancel.trader_id;
1096 let strategy_id = cancel.strategy_id;
1097 let instrument_id = cancel.instrument_id;
1098
1099 self.spawn_task("batch_cancel_order", async move {
1100 if let Err(e) = ws_client
1101 .cancel_order(
1102 &order_id,
1103 client_order_id,
1104 trader_id,
1105 strategy_id,
1106 instrument_id,
1107 )
1108 .await
1109 {
1110 log::error!(
1111 "Batch cancel order failed: order_id={order_id}, client_order_id={client_order_id}, error={e}"
1112 );
1113 anyhow::bail!("Batch cancel order failed: {e}");
1114 }
1115 Ok(())
1116 });
1117 }
1118
1119 Ok(())
1120 }
1121}
1122
1123fn dispatch_ws_message(message: NautilusWsMessage, emitter: &ExecutionEventEmitter) {
1125 match message {
1126 NautilusWsMessage::AccountState(state) => {
1127 emitter.send_account_state(state);
1128 }
1129 NautilusWsMessage::OrderStatusReports(reports) => {
1130 log::debug!("Processing {} order status report(s)", reports.len());
1131 for report in reports {
1132 emitter.send_order_status_report(report);
1133 }
1134 }
1135 NautilusWsMessage::FillReports(reports) => {
1136 log::debug!("Processing {} fill report(s)", reports.len());
1137 for report in reports {
1138 emitter.send_fill_report(report);
1139 }
1140 }
1141 NautilusWsMessage::OrderFilled(event) => {
1142 emitter.send_order_event(OrderEventAny::Filled(event));
1143 }
1144 NautilusWsMessage::OrderRejected(event) => {
1145 emitter.send_order_event(OrderEventAny::Rejected(event));
1146 }
1147 NautilusWsMessage::OrderAccepted(event) => {
1148 emitter.send_order_event(OrderEventAny::Accepted(event));
1149 }
1150 NautilusWsMessage::OrderCanceled(event) => {
1151 emitter.send_order_event(OrderEventAny::Canceled(event));
1152 }
1153 NautilusWsMessage::OrderExpired(event) => {
1154 emitter.send_order_event(OrderEventAny::Expired(event));
1155 }
1156 NautilusWsMessage::OrderUpdated(event) => {
1157 emitter.send_order_event(OrderEventAny::Updated(event));
1158 }
1159 NautilusWsMessage::OrderCancelRejected(event) => {
1160 emitter.send_order_event(OrderEventAny::CancelRejected(event));
1161 }
1162 NautilusWsMessage::OrderModifyRejected(event) => {
1163 emitter.send_order_event(OrderEventAny::ModifyRejected(event));
1164 }
1165 NautilusWsMessage::Error(e) => {
1166 log::warn!("WebSocket error: {e}");
1167 }
1168 NautilusWsMessage::Reconnected => {
1169 log::info!("WebSocket reconnected");
1170 }
1171 NautilusWsMessage::Authenticated(auth) => {
1172 log::debug!("WebSocket authenticated: scope={}", auth.scope);
1173 }
1174 NautilusWsMessage::AuthenticationFailed(reason) => {
1175 log::error!("Authentication failed in execution client: {reason}");
1176 }
1177 NautilusWsMessage::Data(_)
1178 | NautilusWsMessage::Deltas(_)
1179 | NautilusWsMessage::Instrument(_)
1180 | NautilusWsMessage::InstrumentStatus(_)
1181 | NautilusWsMessage::FundingRates(_)
1182 | NautilusWsMessage::OptionGreeks(_)
1183 | NautilusWsMessage::Raw(_) => {
1184 log::trace!("Ignoring data message in execution client");
1186 }
1187 }
1188}
1189
1190fn reject_modify_command(
1191 emitter: &ExecutionEventEmitter,
1192 clock: &AtomicTime,
1193 cmd: &ModifyOrder,
1194 reason: &str,
1195) -> anyhow::Result<()> {
1196 let ts_event = clock.get_time_ns();
1197 emitter.emit_order_modify_rejected_event(
1198 cmd.strategy_id,
1199 cmd.instrument_id,
1200 cmd.client_order_id,
1201 cmd.venue_order_id,
1202 reason,
1203 ts_event,
1204 );
1205 anyhow::bail!("{reason}");
1206}
1207
1208#[cfg(test)]
1209mod tests {
1210 use nautilus_common::messages::{ExecutionEvent, execution::ExecutionReport};
1211 use nautilus_core::UUID4;
1212 use nautilus_model::{
1213 enums::{LiquiditySide, OrderSide},
1214 events::OrderFilled,
1215 identifiers::{ClientOrderId, InstrumentId, StrategyId, TradeId, TraderId, VenueOrderId},
1216 types::{Currency, Money, Price, Quantity},
1217 };
1218 use rstest::rstest;
1219
1220 use super::*;
1221
1222 fn dispatch_test_rig() -> (
1223 ExecutionEventEmitter,
1224 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
1225 ) {
1226 let trader_id = TraderId::from("TRADER-001");
1227 let account_id = AccountId::from("DERIBIT-001");
1228 let mut emitter = ExecutionEventEmitter::new(
1229 get_atomic_clock_realtime(),
1230 trader_id,
1231 account_id,
1232 AccountType::Margin,
1233 None,
1234 );
1235 let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
1236 emitter.set_sender(tx);
1237 (emitter, rx)
1238 }
1239
1240 fn fill_report() -> FillReport {
1241 FillReport::new(
1242 AccountId::from("DERIBIT-001"),
1243 InstrumentId::from("BTC-PERPETUAL.DERIBIT"),
1244 VenueOrderId::from("ETH-584830574"),
1245 TradeId::from("ETH-2696068"),
1246 OrderSide::Buy,
1247 Quantity::from("1.000000"),
1248 Price::from("203.80"),
1249 Money::from("0.00036801 USDT"),
1250 LiquiditySide::Taker,
1251 Some(ClientOrderId::from("O-19700101-000000-001-001-1")),
1252 None,
1253 UnixNanos::from(2),
1254 UnixNanos::from(3),
1255 None,
1256 )
1257 }
1258
1259 #[rstest]
1260 fn dispatch_tracked_fill_uses_order_event_path() {
1261 let (emitter, mut rx) = dispatch_test_rig();
1262 let report = fill_report();
1263 let filled = OrderFilled::new(
1264 TraderId::from("TRADER-001"),
1265 StrategyId::from("S-001"),
1266 report.instrument_id,
1267 report.client_order_id.unwrap(),
1268 report.venue_order_id,
1269 report.account_id,
1270 report.trade_id,
1271 report.order_side,
1272 OrderType::Market,
1273 report.last_qty,
1274 report.last_px,
1275 Currency::USDT(),
1276 report.liquidity_side,
1277 UUID4::new(),
1278 report.ts_event,
1279 report.ts_init,
1280 false,
1281 None,
1282 Some(report.commission),
1283 None,
1284 );
1285
1286 dispatch_ws_message(NautilusWsMessage::OrderFilled(filled), &emitter);
1287
1288 assert!(matches!(
1289 rx.try_recv().unwrap(),
1290 ExecutionEvent::Order(OrderEventAny::Filled(event))
1291 if event.client_order_id == ClientOrderId::from("O-19700101-000000-001-001-1")
1292 && event.trade_id == TradeId::from("ETH-2696068")
1293 ));
1294 assert!(rx.try_recv().is_err());
1295 }
1296
1297 #[rstest]
1298 fn dispatch_untracked_fill_keeps_report_path() {
1299 let (emitter, mut rx) = dispatch_test_rig();
1300 let report = fill_report();
1301
1302 dispatch_ws_message(NautilusWsMessage::FillReports(vec![report]), &emitter);
1303
1304 assert!(matches!(
1305 rx.try_recv().unwrap(),
1306 ExecutionEvent::Report(ExecutionReport::Fill(_))
1307 ));
1308 assert!(rx.try_recv().is_err());
1309 }
1310}