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