Skip to main content

nautilus_deribit/
execution.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution client implementation for the Deribit adapter.
17
18use 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/// Deribit live execution client.
68#[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    /// Creates a new [`DeribitExecutionClient`].
82    ///
83    /// # Errors
84    ///
85    /// Returns an error if the client fails to initialize.
86    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        // Set account ID for order/fill reports
143        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    /// Spawns an async task for execution operations.
170    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    /// Aborts all pending async tasks.
186    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    // Rejects unsupported order types and time-in-force values
241    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            // Deribit rejects `time_in_force` on market-style order types
260            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        // Deribit's `valid_until` is a REQUEST timeout, not order expiry.
285        // Deribit's `good_til_day` expires at end of trading session (8 UTC).
286        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    /// Submits a single order to Deribit.
320    ///
321    /// This is the core submission logic shared by `submit_order` and `submit_order_list`.
322    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    /// Spawns a stream handler to dispatch WebSocket messages to the execution engine.
378    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        // Check if credentials are available before requesting account state
507        if !self.config.has_api_credentials() {
508            anyhow::bail!("Missing API credentials; set Deribit environment variables");
509        }
510
511        // Set account ID for order/fill reports
512        self.ws_client.set_account_id(self.core.account_id);
513
514        // Fetch and cache instruments in both HTTP client and WebSocket client
515        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        // Fetch initial account state
538        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            // Subscribe to user order and trade updates for all instruments
560            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                // Roll back subscription state so a retry re-sends subscribe requests
575                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            // Spawn stream handler to dispatch WebSocket messages to the execution engine
584            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 venue_order_id is provided, fetch the specific order by ID
617        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 client_order_id is provided, search open then closed orders
652        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, // search all orders, not just open
661                )
662                .await?;
663
664            // Filter by client_order_id
665            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        // Filter by venue_order_id if provided
700        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) // get all orders for mass status
730            .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        // Extract venue order ID (Deribit's order_id)
793        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        // Spawn async task to query order state via WebSocket
807        // Response will be dispatched through the WebSocket stream handler as OrderStatusReport
808        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        // Deribit doesn't have native batch order submission
847        // Loop through and submit each order individually using shared logic
848        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        // Extract venue order ID (Deribit's order_id)
859        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        // Extract quantity - if not provided, get from order in cache
872        let quantity = if let Some(qty) = cmd.quantity {
873            qty
874        } else {
875            // Get order from cache to use its current quantity
876            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        // Spawn async task to send modify via WebSocket
912        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        // Extract venue order ID (Deribit's order_id)
940        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        // Spawn async task to send cancel via WebSocket
959        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        // Without a side filter, use efficient bulk cancel via Deribit API
985        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        // For specific side (Buy/Sell), filter from cache and cancel individually
1003        // Deribit API doesn't support side filtering, so we implement it locally
1004        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        // Cancel each matching order individually
1038        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        // Deribit doesn't have native batch cancel by order ID
1080        // Loop through and cancel each order individually
1081        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
1123/// Dispatches a WebSocket message using the event emitter.
1124fn 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            // Data messages are handled by the data client, not execution
1185            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}