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    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/// Deribit live execution client.
69#[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    /// Creates a new [`DeribitExecutionClient`].
83    ///
84    /// # Errors
85    ///
86    /// Returns an error if the client fails to initialize.
87    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        // Set account ID for order/fill reports
132        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    /// Spawns an async task for execution operations.
159    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    /// Aborts all pending async tasks.
175    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    // Rejects unsupported order types and time-in-force values
230    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            // Deribit rejects `time_in_force` on market-style order types
249            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        // Deribit's `valid_until` is a REQUEST timeout, not order expiry.
274        // Deribit's `good_til_day` expires at end of trading session (8 UTC).
275        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    /// Submits a single order to Deribit.
309    ///
310    /// This is the core submission logic shared by `submit_order` and `submit_order_list`.
311    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    /// Spawns a stream handler to dispatch WebSocket messages to the execution engine.
367    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        // Check if credentials are available before requesting account state
496        if !self.config.has_api_credentials() {
497            anyhow::bail!("Missing API credentials; set Deribit environment variables");
498        }
499
500        // Set account ID for order/fill reports
501        self.ws_client.set_account_id(self.core.account_id);
502
503        // Fetch and cache instruments in both HTTP client and WebSocket client
504        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        // Fetch initial account state
527        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            // Subscribe to user order and trade updates for all instruments
549            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                // Roll back subscription state so a retry re-sends subscribe requests
564                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            // Spawn stream handler to dispatch WebSocket messages to the execution engine
573            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 venue_order_id is provided, fetch the specific order by ID
606        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 client_order_id is provided, search open then closed orders
641        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, // search all orders, not just open
650                )
651                .await?;
652
653            // Filter by client_order_id
654            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        // Filter by venue_order_id if provided
689        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) // get all orders for mass status
721            .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        // Extract venue order ID (Deribit's order_id)
784        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        // Spawn async task to query order state via WebSocket
798        // Response will be dispatched through the WebSocket stream handler as OrderStatusReport
799        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        // Deribit doesn't have native batch order submission
838        // Loop through and submit each order individually using shared logic
839        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        // Extract venue order ID (Deribit's order_id)
850        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        // Extract quantity - if not provided, get from order in cache
863        let quantity = if let Some(qty) = cmd.quantity {
864            qty
865        } else {
866            // Get order from cache to use its current quantity
867            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        // Spawn async task to send modify via WebSocket
903        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        // Extract venue order ID (Deribit's order_id)
931        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        // Spawn async task to send cancel via WebSocket
950        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        // Without a side filter, use efficient bulk cancel via Deribit API
976        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        // For specific side (Buy/Sell), filter from cache and cancel individually
994        // Deribit API doesn't support side filtering, so we implement it locally
995        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        // Cancel each matching order individually
1029        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        // Deribit doesn't have native batch cancel by order ID
1071        // Loop through and cancel each order individually
1072        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
1114/// Dispatches a WebSocket message using the event emitter.
1115fn 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            // Data messages are handled by the data client, not execution
1176            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}