1use std::{
42 sync::{
43 Arc,
44 atomic::{AtomicU64, Ordering},
45 },
46 time::{Duration, Instant},
47};
48
49use ahash::AHashMap;
50use anyhow::Context;
51use async_trait::async_trait;
52use dashmap::DashMap;
53use futures_util::{Stream, StreamExt, pin_mut};
54use nautilus_common::{
55 clients::ExecutionClient,
56 live::runner::get_exec_event_sender,
57 messages::execution::{
58 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
59 GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
60 ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
61 },
62};
63use nautilus_core::{
64 Params, UUID4, UnixNanos,
65 time::{AtomicTime, get_atomic_clock_realtime},
66};
67use nautilus_live::{
68 ExecutionClientCore, ExecutionEventEmitter, SocketControlFactory,
69 task::{TaskGroup, TaskGroupGuard},
70};
71use nautilus_model::{
72 accounts::AccountAny,
73 enums::{AccountType, OmsType, OrderStatus, OrderType, TimeInForce},
74 events::{AccountState, OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired},
75 identifiers::{
76 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Symbol, Venue, VenueOrderId,
77 },
78 instruments::{Instrument, InstrumentAny},
79 orders::{Order, OrderAny},
80 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
81 types::{AccountBalance, Currency, MarginBalance, Money},
82};
83use nautilus_network::retry::RetryConfig;
84use parking_lot::Mutex;
85use rust_decimal::Decimal;
86
87use crate::{
88 common::{
89 consts::DYDX_VENUE,
90 credential::{DydxCredential, credential_env_vars},
91 instrument_cache::InstrumentCache,
92 parse::nanos_to_secs_i64,
93 },
94 config::DydxAdapterConfig,
95 error::DydxError,
96 execution::{
97 broadcaster::TxBroadcaster,
98 encoder::ClientOrderIdEncoder,
99 order_builder::OrderMessageBuilder,
100 tx_manager::TransactionManager,
101 types::{LimitOrderParams, OrderContext},
102 },
103 grpc::{DydxGrpcClient, SHORT_TERM_ORDER_MAXIMUM_LIFETIME, types::ChainId},
104 http::{
105 client::DydxHttpClient,
106 parse::{
107 parse_account_state, parse_fill_report, parse_order_status_report,
108 parse_position_status_report,
109 },
110 },
111 websocket::{
112 DydxWsDispatchState, OrderIdentity,
113 client::DydxWebSocketClient,
114 enums::DydxWsOutputMessage,
115 fill_report_to_order_filled,
116 parse::{parse_ws_fill_report, parse_ws_order_report, parse_ws_position_report},
117 },
118};
119
120pub mod block_time;
121pub mod broadcaster;
122pub mod encoder;
123pub mod order_builder;
124pub mod submitter;
125pub mod tx_manager;
126pub mod types;
127pub mod wallet;
128
129use block_time::BlockTimeMonitor;
130
131const DYDX_INDEXER_REPORT_LIMIT: u32 = 1_000;
132
133type CancelAllOrderData = (
134 StrategyId,
135 ClientOrderId,
136 Option<VenueOrderId>,
137 TimeInForce,
138 Option<UnixNanos>,
139);
140
141#[derive(Clone, Copy)]
142struct DydxCancelOrderRequest {
143 instrument_id: InstrumentId,
144 client_id: u32,
145 order_flags: u32,
146 strategy_id: StrategyId,
147 client_order_id: ClientOrderId,
148 venue_order_id: Option<VenueOrderId>,
149}
150
151fn apply_avg_px_from_fills(order_reports: &mut [OrderStatusReport], fill_reports: &[FillReport]) {
152 let mut totals: AHashMap<VenueOrderId, (Decimal, Decimal)> = AHashMap::new();
153 for fill in fill_reports {
154 let entry = totals.entry(fill.venue_order_id).or_default();
155 let qty = fill.last_qty.as_decimal();
156 entry.0 += fill.last_px.as_decimal() * qty;
157 entry.1 += qty;
158 }
159
160 for report in order_reports {
161 if let Some((notional, total_qty)) = totals.get(&report.venue_order_id)
162 && !total_qty.is_zero()
163 {
164 report.avg_px = Some(notional / total_qty);
165 }
166 }
167}
168
169#[derive(Debug)]
184pub struct DydxExecutionClient {
185 core: ExecutionClientCore,
186 clock: &'static AtomicTime,
187 config: DydxAdapterConfig,
188 emitter: ExecutionEventEmitter,
189 http_client: DydxHttpClient,
190 ws_client: DydxWebSocketClient,
191 grpc_client: Arc<tokio::sync::RwLock<Option<DydxGrpcClient>>>,
192 instrument_cache: Arc<InstrumentCache>,
193 block_time_monitor: Arc<BlockTimeMonitor>,
194 oracle_prices: Arc<DashMap<InstrumentId, Decimal>>,
195 encoder: Arc<ClientOrderIdEncoder>,
196 dispatch_state: Arc<DydxWsDispatchState>,
197 order_contexts: Arc<DashMap<u32, OrderContext>>,
198 order_id_map: Arc<DashMap<String, (u32, u32)>>,
199 wallet_address: String,
200 subaccount_number: u32,
201 tx_manager: Option<Arc<TransactionManager>>,
202 broadcaster: Option<Arc<TxBroadcaster>>,
203 order_builder: Option<Arc<OrderMessageBuilder>>,
204 session_tasks: TaskGroup,
205 pending_tasks: TaskGroup,
206 shutdown_errors: Vec<String>,
207 pending_task_labels: Arc<Mutex<AHashMap<u64, &'static str>>>,
208 next_pending_task_id: AtomicU64,
209}
210
211impl DydxExecutionClient {
212 pub fn new(
218 core: ExecutionClientCore,
219 config: DydxAdapterConfig,
220 wallet_address: String,
221 subaccount_number: u32,
222 ) -> anyhow::Result<Self> {
223 let trader_id = core.trader_id;
224 let account_id = core.account_id;
225 let clock = get_atomic_clock_realtime();
226 let emitter =
227 ExecutionEventEmitter::new(clock, trader_id, account_id, AccountType::Margin, None);
228
229 let retry_config = RetryConfig {
230 max_retries: config.max_retries,
231 initial_delay_ms: config.retry_delay_initial_ms,
232 max_delay_ms: config.retry_delay_max_ms,
233 ..Default::default()
234 };
235 let http_client = DydxHttpClient::new(
236 Some(config.base_url.clone()),
237 config.timeout_secs,
238 config.proxy_url.clone(),
239 config.network,
240 Some(retry_config),
241 )?;
242
243 let instrument_cache = http_client.instrument_cache().clone();
245
246 let credential = DydxCredential::resolve(
248 config.private_key.as_deref(),
249 config.network,
250 config.authenticator_ids.clone(),
251 )?
252 .ok_or_else(|| anyhow::anyhow!("Credentials required for execution client"))?;
253
254 let ws_client = DydxWebSocketClient::new_private_with_cache(
256 config.ws_url.clone(),
257 credential,
258 core.account_id,
259 instrument_cache.clone(),
260 Some(20),
261 config.transport_backend,
262 config.proxy_url.clone(),
263 )
264 .with_socket_factory(SocketControlFactory::new(core.client_id, Some(*DYDX_VENUE)));
265
266 let grpc_client = Arc::new(tokio::sync::RwLock::new(None));
267
268 let session_tasks = TaskGroup::new();
269 let pending_tasks = TaskGroup::new();
270
271 Ok(Self {
272 core,
273 clock,
274 config,
275 emitter,
276 http_client,
277 ws_client,
278 grpc_client,
279 instrument_cache,
280 block_time_monitor: Arc::new(BlockTimeMonitor::new()),
281 oracle_prices: Arc::new(DashMap::new()),
282 encoder: Arc::new(ClientOrderIdEncoder::new()),
283 dispatch_state: Arc::new(DydxWsDispatchState::default()),
284 order_contexts: Arc::new(DashMap::new()),
285 order_id_map: Arc::new(DashMap::new()),
286 wallet_address,
287 subaccount_number,
288 tx_manager: None,
289 broadcaster: None,
290 order_builder: None,
291 session_tasks,
292 pending_tasks,
293 shutdown_errors: Vec::new(),
294 pending_task_labels: Arc::new(Mutex::new(AHashMap::new())),
295 next_pending_task_id: AtomicU64::new(0),
296 })
297 }
298
299 fn resolve_private_key(config: &DydxAdapterConfig) -> anyhow::Result<String> {
300 let (private_key_env, _) = credential_env_vars(config.network);
301
302 if let Some(ref pk) = config.private_key
304 && !pk.trim().is_empty()
305 {
306 return Ok(pk.clone());
307 }
308
309 if let Some(pk) = std::env::var(private_key_env)
311 .ok()
312 .filter(|s| !s.trim().is_empty())
313 {
314 return Ok(pk);
315 }
316
317 anyhow::bail!("{private_key_env} not found in config or environment")
318 }
319
320 fn register_order_context(&self, client_id_u32: u32, context: OrderContext) {
321 self.order_contexts.insert(client_id_u32, context);
322 }
323
324 fn get_order_context(&self, client_id_u32: u32) -> Option<OrderContext> {
325 self.order_contexts
326 .get(&client_id_u32)
327 .map(|r| r.value().clone())
328 }
329
330 fn get_chain_id(&self) -> ChainId {
331 self.config.get_chain_id()
332 }
333
334 fn spawn_ws_stream_handler(
335 &self,
336 stream: impl Stream<Item = DydxWsOutputMessage> + Send + 'static,
337 ) -> anyhow::Result<()> {
338 if !self.session_tasks.is_empty() {
339 anyhow::bail!("dYdX WebSocket stream task is already registered");
340 }
341
342 log::debug!("Starting execution WebSocket message processing task");
343
344 let trader_id = self.core.trader_id;
346 let account_id = self.core.account_id;
347 let instrument_cache = self.instrument_cache.clone();
348 let oracle_prices = self.oracle_prices.clone();
349 let encoder = self.encoder.clone();
350 let order_contexts = self.order_contexts.clone();
351 let order_id_map = self.order_id_map.clone();
352 let dispatch_state = self.dispatch_state.clone();
353 let block_time_monitor = self.block_time_monitor.clone();
354 let emitter = self.emitter.clone();
355 let clock = self.clock;
356
357 let future = async move {
358 log::debug!("Execution WebSocket message loop started");
359
360 let mut cum_fill_totals: AHashMap<VenueOrderId, (Decimal, Decimal)> = AHashMap::new();
362
363 pin_mut!(stream);
364 while let Some(msg) = stream.next().await {
365 match msg {
366 DydxWsOutputMessage::SubaccountSubscribed(msg) => {
367 log::debug!("Parsing subaccount subscription with full context");
368
369 let inst_map = instrument_cache.to_instrument_id_map();
370
371 let oracle_map: std::collections::HashMap<_, _> = oracle_prices
372 .iter()
373 .map(|entry| (*entry.key(), *entry.value()))
374 .collect();
375
376 let ts_init = clock.get_time_ns();
377 let ts_event = ts_init;
378
379 if let Some(ref subaccount) = msg.contents.subaccount {
380 match parse_account_state(
381 subaccount,
382 account_id,
383 &inst_map,
384 &oracle_map,
385 ts_event,
386 ts_init,
387 ) {
388 Ok(account_state) => {
389 log::debug!(
390 "Parsed account state: {} balance(s), {} margin(s)",
391 account_state.balances.len(),
392 account_state.margins.len()
393 );
394 emitter.send_account_state(account_state);
395 }
396 Err(e) => {
397 log::error!("Failed to parse account state: {e}");
398 }
399 }
400
401 if let Some(ref positions) = subaccount.open_perpetual_positions {
402 log::debug!(
403 "Parsing {} position(s) from subscription",
404 positions.len()
405 );
406
407 for (market, ws_position) in positions {
408 match parse_ws_position_report(
409 ws_position,
410 &instrument_cache,
411 account_id,
412 ts_init,
413 ) {
414 Ok(report) => {
415 log::debug!(
416 "Parsed position report: {} {} {} {}",
417 report.instrument_id,
418 report.position_side,
419 report.quantity,
420 market
421 );
422 emitter.send_position_report(report);
423 }
424 Err(e) => {
425 log::error!(
426 "Failed to parse WebSocket position for {market}: {e}"
427 );
428 }
429 }
430 }
431 }
432 } else {
433 log::warn!(
434 "Subaccount subscription without initial state (new/empty subaccount)"
435 );
436
437 let currency =
438 Currency::get_or_create_crypto_with_context("USDC", None);
439 let zero = Money::zero(currency);
440 let balance = AccountBalance::new_checked(zero, zero, zero)
441 .expect("zero balance should always be valid");
442 let account_state = AccountState::new(
443 account_id,
444 AccountType::Margin,
445 vec![balance],
446 vec![],
447 true,
448 UUID4::new(),
449 ts_init,
450 ts_init,
451 None,
452 );
453 emitter.send_account_state(account_state);
454 }
455 }
456 DydxWsOutputMessage::SubaccountsChannelData(data) => {
457 log::debug!(
458 "Processing subaccounts channel data (orders={:?}, fills={:?})",
459 data.contents.orders.as_ref().map(|o| o.len()),
460 data.contents.fills.as_ref().map(|f| f.len())
461 );
462 let ts_init = clock.get_time_ns();
463
464 let mut terminal_orders: Vec<(u32, u32, String)> = Vec::new();
465 let mut pending_order_reports = Vec::new();
466
467 if let Some(ref orders) = data.contents.orders {
469 for ws_order in orders {
470 log::debug!(
471 "Parsing WS order: clob_pair_id={}, status={:?}, client_id={}",
472 ws_order.clob_pair_id,
473 ws_order.status,
474 ws_order.client_id
475 );
476
477 if let Ok(client_id_u32) = ws_order.client_id.parse::<u32>() {
478 let client_meta = ws_order
479 .client_metadata
480 .as_ref()
481 .and_then(|s| s.parse::<u32>().ok())
482 .unwrap_or(crate::grpc::DEFAULT_RUST_CLIENT_METADATA);
483 order_id_map
484 .insert(ws_order.id.clone(), (client_id_u32, client_meta));
485 }
486
487 match parse_ws_order_report(
488 ws_order,
489 &instrument_cache,
490 &order_contexts,
491 &encoder,
492 account_id,
493 ts_init,
494 ) {
495 Ok(report) => {
496 if !report.order_status.is_open()
497 && let Ok(cid) = ws_order.client_id.parse::<u32>()
498 {
499 let meta = ws_order
500 .client_metadata
501 .as_ref()
502 .and_then(|s| s.parse::<u32>().ok())
503 .unwrap_or(
504 crate::grpc::DEFAULT_RUST_CLIENT_METADATA,
505 );
506 terminal_orders.push((cid, meta, ws_order.id.clone()));
507 }
508 let order_side = report
509 .order_side
510 .as_ref()
511 .map_or("NO_ORDER_SIDE", AsRef::as_ref);
512 log::debug!(
513 "Parsed order report: {} {} {:?} qty={} client_order_id={:?}",
514 report.instrument_id,
515 order_side,
516 report.order_status,
517 report.quantity,
518 report.client_order_id
519 );
520 pending_order_reports.push(report);
521 }
522 Err(e) => {
523 log::error!("Failed to parse WebSocket order: {e}");
524 }
525 }
526 }
527 }
528
529 if let Some(ref fills) = data.contents.fills {
531 for ws_fill in fills {
532 match parse_ws_fill_report(
533 ws_fill,
534 &instrument_cache,
535 &order_id_map,
536 &order_contexts,
537 &encoder,
538 account_id,
539 ts_init,
540 ) {
541 Ok(report) => {
542 log::debug!(
543 "Parsed fill report: {} {} {} @ {} client_order_id={:?}",
544 report.instrument_id,
545 report.venue_order_id,
546 report.last_qty,
547 report.last_px,
548 report.client_order_id
549 );
550
551 let identity = report.client_order_id.and_then(|cid| {
552 dispatch_state
553 .order_identities
554 .get(&cid)
555 .map(|r| (cid, r.clone()))
556 });
557
558 if let Some((cid, ident)) = identity {
559 if !dispatch_state.emitted_accepted.contains(&cid) {
561 dispatch_state.insert_accepted(cid);
562 let accepted = OrderAccepted::new(
563 trader_id,
564 ident.strategy_id,
565 ident.instrument_id,
566 cid,
567 report.venue_order_id,
568 account_id,
569 UUID4::new(),
570 ts_init,
571 ts_init,
572 false,
573 );
574 emitter.send_order_event(OrderEventAny::Accepted(
575 accepted,
576 ));
577 }
578
579 dispatch_state.insert_filled(cid);
580 let instrument =
581 instrument_cache.get(&report.instrument_id);
582 let quote_currency = instrument
583 .map_or_else(Currency::USD, |i: InstrumentAny| {
584 i.quote_currency()
585 });
586 let filled = fill_report_to_order_filled(
587 &report,
588 trader_id,
589 &ident,
590 quote_currency,
591 );
592 emitter.send_order_event(OrderEventAny::Filled(filled));
593 } else {
594 let entry = cum_fill_totals
596 .entry(report.venue_order_id)
597 .or_default();
598 let qty = report.last_qty.as_decimal();
599 entry.0 += report.last_px.as_decimal() * qty;
600 entry.1 += qty;
601 emitter.send_fill_report(report);
602 }
603 }
604 Err(e) => {
605 log::error!("Failed to parse WebSocket fill: {e}");
606 }
607 }
608 }
609 }
610
611 for report in &mut pending_order_reports {
614 if let Some((notional, total_qty)) =
615 cum_fill_totals.get(&report.venue_order_id)
616 && !total_qty.is_zero()
617 {
618 report.avg_px = Some(notional / total_qty);
619 }
620 }
621
622 for report in pending_order_reports {
623 let identity = report.client_order_id.and_then(|cid| {
624 dispatch_state
625 .order_identities
626 .get(&cid)
627 .map(|r| (cid, r.clone()))
628 });
629
630 if let Some((cid, ident)) = identity {
631 match report.order_status {
633 OrderStatus::Accepted => {
634 if dispatch_state.emitted_accepted.contains(&cid)
635 || dispatch_state.filled_orders.contains(&cid)
636 {
637 log::debug!("Skipping duplicate Accepted for {cid}");
638 continue;
639 }
640 dispatch_state.insert_accepted(cid);
641 let accepted = OrderAccepted::new(
642 trader_id,
643 ident.strategy_id,
644 ident.instrument_id,
645 cid,
646 report.venue_order_id,
647 account_id,
648 UUID4::new(),
649 report.ts_last,
650 ts_init,
651 false,
652 );
653 emitter.send_order_event(OrderEventAny::Accepted(accepted));
654 }
655 OrderStatus::Canceled => {
656 if !dispatch_state.emitted_accepted.contains(&cid) {
658 dispatch_state.insert_accepted(cid);
659 let accepted = OrderAccepted::new(
660 trader_id,
661 ident.strategy_id,
662 ident.instrument_id,
663 cid,
664 report.venue_order_id,
665 account_id,
666 UUID4::new(),
667 ts_init,
668 ts_init,
669 false,
670 );
671 emitter.send_order_event(OrderEventAny::Accepted(
672 accepted,
673 ));
674 }
675 let is_expiry = report
678 .expire_time
679 .is_some_and(|exp| report.ts_last >= exp);
680
681 if is_expiry {
682 let expired = OrderExpired::new(
683 trader_id,
684 ident.strategy_id,
685 ident.instrument_id,
686 cid,
687 UUID4::new(),
688 report.ts_last,
689 ts_init,
690 false,
691 Some(report.venue_order_id),
692 Some(account_id),
693 );
694 emitter
695 .send_order_event(OrderEventAny::Expired(expired));
696 } else {
697 let canceled = OrderCanceled::new(
698 trader_id,
699 ident.strategy_id,
700 ident.instrument_id,
701 cid,
702 UUID4::new(),
703 report.ts_last,
704 ts_init,
705 false,
706 Some(report.venue_order_id),
707 Some(account_id),
708 );
709 emitter.send_order_event(OrderEventAny::Canceled(
710 canceled,
711 ));
712 }
713 dispatch_state.cleanup_terminal(&cid);
714 }
715 OrderStatus::Filled => {
716 dispatch_state.cleanup_terminal(&cid);
718 }
719 OrderStatus::Expired => {
720 if !dispatch_state.emitted_accepted.contains(&cid) {
726 dispatch_state.insert_accepted(cid);
727 let accepted = OrderAccepted::new(
728 trader_id,
729 ident.strategy_id,
730 ident.instrument_id,
731 cid,
732 report.venue_order_id,
733 account_id,
734 UUID4::new(),
735 ts_init,
736 ts_init,
737 false,
738 );
739 emitter.send_order_event(OrderEventAny::Accepted(
740 accepted,
741 ));
742 }
743 let expired = OrderExpired::new(
744 trader_id,
745 ident.strategy_id,
746 ident.instrument_id,
747 cid,
748 UUID4::new(),
749 report.ts_last,
750 ts_init,
751 false,
752 Some(report.venue_order_id),
753 Some(account_id),
754 );
755 emitter.send_order_event(OrderEventAny::Expired(expired));
756 dispatch_state.cleanup_terminal(&cid);
757 }
758 _ => {
759 emitter.send_order_status_report(report);
761 }
762 }
763 } else {
764 emitter.send_order_status_report(report);
766 }
767 }
768
769 for (client_id, client_metadata, order_id) in terminal_orders {
771 order_contexts.remove(&client_id);
772 encoder.remove(client_id, client_metadata);
773 order_id_map.remove(&order_id);
774 cum_fill_totals.remove(&VenueOrderId::new(&order_id));
775 }
776 }
777 DydxWsOutputMessage::Markets(contents) => {
778 if let Some(ref oracle_map) = contents.oracle_prices {
779 for (symbol_str, oracle_market) in oracle_map {
780 let instrument_id = {
781 let symbol = format!("{symbol_str}-PERP");
782 InstrumentId::new(
783 Symbol::new(&symbol),
784 *crate::common::consts::DYDX_VENUE,
785 )
786 };
787
788 if instrument_cache.get(&instrument_id).is_some()
789 && let Ok(price_dec) =
790 oracle_market.oracle_price.parse::<Decimal>()
791 {
792 oracle_prices.insert(instrument_id, price_dec);
793 log::trace!(
794 "Updated oracle price for {instrument_id}: {price_dec}"
795 );
796 }
797 }
798 }
799
800 if let Some(ref markets) = contents.markets {
801 for (symbol_str, market_data) in markets {
802 if let Some(oracle_price_str) = &market_data.oracle_price {
803 let instrument_id = {
804 let symbol = format!("{symbol_str}-PERP");
805 InstrumentId::new(
806 Symbol::new(&symbol),
807 *crate::common::consts::DYDX_VENUE,
808 )
809 };
810
811 if instrument_cache.get(&instrument_id).is_some()
812 && let Ok(price_dec) = oracle_price_str.parse::<Decimal>()
813 {
814 oracle_prices.insert(instrument_id, price_dec);
815 }
816 }
817 }
818 }
819 }
820 DydxWsOutputMessage::BlockHeight { height, time } => {
821 log::debug!("Block height update: {height} at {time}");
822 block_time_monitor.record_block(height, time);
823 }
824 DydxWsOutputMessage::Error(err) => {
825 log::warn!("WebSocket error: {err:?}");
826 }
827 DydxWsOutputMessage::Reconnected { .. } => {
828 log::info!("WebSocket reconnected");
829 }
830 _ => {}
831 }
832 }
833 log::debug!("WebSocket message processing task ended");
834 };
835
836 self.session_tasks
837 .spawn(future)
838 .context("failed to register dYdX execution WebSocket stream task")?;
839 log::debug!("WebSocket stream handler started");
840 Ok(())
841 }
842
843 fn mark_instruments_initialized(&self) {
848 let count = self.instrument_cache.len();
849 self.core.set_instruments_initialized();
850 log::debug!("Instruments initialized: {count} instruments in shared cache");
851 }
852
853 fn get_instrument_by_market(&self, market: &str) -> Option<InstrumentAny> {
854 self.instrument_cache.get_by_market(market)
855 }
856
857 fn get_instrument_by_clob_pair_id(&self, clob_pair_id: u32) -> Option<InstrumentAny> {
858 let instrument = self.instrument_cache.get_by_clob_id(clob_pair_id);
859
860 if instrument.is_none() {
861 self.instrument_cache.log_missing_clob_pair_id(clob_pair_id);
862 }
863
864 instrument
865 }
866
867 fn get_execution_components(
871 &self,
872 ) -> anyhow::Result<(
873 Arc<TransactionManager>,
874 Arc<TxBroadcaster>,
875 Arc<OrderMessageBuilder>,
876 )> {
877 let tx_manager = self
878 .tx_manager
879 .as_ref()
880 .ok_or_else(|| {
881 anyhow::anyhow!("TransactionManager not initialized - call connect() first")
882 })?
883 .clone();
884 let broadcaster = self
885 .broadcaster
886 .as_ref()
887 .ok_or_else(|| anyhow::anyhow!("TxBroadcaster not initialized - call connect() first"))?
888 .clone();
889 let order_builder = self
890 .order_builder
891 .as_ref()
892 .ok_or_else(|| {
893 anyhow::anyhow!("OrderMessageBuilder not initialized - call connect() first")
894 })?
895 .clone();
896 Ok((tx_manager, broadcaster, order_builder))
897 }
898
899 fn spawn_task<F>(&self, label: &'static str, fut: F)
900 where
901 F: Future<Output = anyhow::Result<()>> + Send + 'static,
902 {
903 let future = async move {
904 if let Err(e) = fut.await {
905 log::error!("{label}: {e:?}");
906 }
907 };
908
909 self.spawn_labeled(label, future);
910 }
911
912 fn spawn_order_task<F>(
918 &self,
919 label: &'static str,
920 strategy_id: StrategyId,
921 instrument_id: InstrumentId,
922 client_order_id: ClientOrderId,
923 fut: F,
924 ) where
925 F: Future<Output = anyhow::Result<()>> + Send + 'static,
926 {
927 let emitter = self.emitter.clone();
928 let clock = self.clock;
929
930 let future = async move {
931 if let Err(e) = fut.await {
932 if is_definitive_broadcast_rejection(&e) {
933 let error_msg = format!("{label} failed: {e:?}");
934 log::error!("{error_msg}");
935
936 let ts_event = clock.get_time_ns();
937 emitter.emit_order_rejected_event(
938 strategy_id,
939 instrument_id,
940 client_order_id,
941 &error_msg,
942 ts_event,
943 false,
944 );
945 } else {
946 log::warn!(
947 "Ambiguous dYdX {label} failure for {client_order_id}, awaiting reconciliation: {e:?}"
948 );
949 }
950 }
951 };
952
953 self.spawn_labeled(label, future);
954 }
955
956 fn spawn_labeled<F>(&self, label: &'static str, future: F)
957 where
958 F: Future<Output = ()> + Send + 'static,
959 {
960 let id = self.next_pending_task_id.fetch_add(1, Ordering::Relaxed);
961 self.pending_task_labels.lock().insert(id, label);
962
963 let label_guard = PendingTaskLabel {
964 id,
965 labels: Arc::clone(&self.pending_task_labels),
966 };
967 let wrapped = async move {
968 let _label = label_guard;
969 future.await;
970 };
971
972 if let Err(e) = self.pending_tasks.spawn(wrapped) {
973 self.pending_task_labels.lock().remove(&id);
974 log::warn!("Skipping dYdX {label} after shutdown began: {e}");
975 }
976 }
977
978 fn begin_pending_shutdown(&self) {
979 let labels = self.pending_task_labels.lock();
980 let pending_cancels = labels
981 .values()
982 .filter(|label| label.contains("cancel"))
983 .count();
984 let pending_other = labels.len() - pending_cancels;
985
986 if pending_cancels > 0 {
987 log::warn!(
988 "Waiting for {pending_cancels} in-flight cancel task(s) before disconnect; orders may still be open on the venue"
989 );
990 }
991
992 if pending_other > 0 {
993 log::debug!("Waiting for {pending_other} other in-flight task(s) before disconnect");
994 }
995 drop(labels);
996
997 self.pending_tasks.begin_shutdown();
998 }
999
1000 async fn finish_tasks(&self) -> anyhow::Result<()> {
1001 let (session_result, pending_result) = tokio::join!(
1002 self.session_tasks
1003 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
1004 self.pending_tasks
1005 .finish_shutdown(Duration::from_secs(2), Duration::from_secs(2)),
1006 );
1007 session_result.context("failed to finish dYdX execution session tasks")?;
1008 pending_result.context("failed to finish dYdX execution command tasks")?;
1009 Ok(())
1010 }
1011
1012 async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
1013 if !self.session_tasks.is_open() || !self.pending_tasks.is_open() {
1014 self.session_tasks.begin_shutdown();
1015 self.begin_pending_shutdown();
1016 self.ws_client.begin_shutdown();
1017 self.finish_shutdown().await?;
1018 self.session_tasks
1019 .start_generation()
1020 .context("failed to start dYdX execution session task generation")?;
1021 self.pending_tasks
1022 .start_generation()
1023 .context("failed to start dYdX execution command task generation")?;
1024 }
1025 Ok(())
1026 }
1027
1028 async fn finish_shutdown(&mut self) -> anyhow::Result<()> {
1029 if let Err(e) = self
1030 .ws_client
1031 .disconnect()
1032 .await
1033 .context("failed to disconnect dYdX websocket")
1034 {
1035 self.shutdown_errors.push(e.to_string());
1036 }
1037
1038 if let Err(e) = self.finish_tasks().await {
1039 self.shutdown_errors.push(e.to_string());
1040 }
1041
1042 if !self.shutdown_errors.is_empty() {
1043 anyhow::bail!(std::mem::take(&mut self.shutdown_errors).join("; "));
1044 }
1045 Ok(())
1046 }
1047
1048 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
1049 self.begin_pending_shutdown();
1050 self.session_tasks.begin_shutdown();
1051 self.ws_client.begin_shutdown();
1052 let shutdown_result = self.finish_shutdown().await;
1053 self.core.set_disconnected();
1054 shutdown_result
1055 }
1056
1057 fn send_modify_rejected(
1059 &self,
1060 strategy_id: StrategyId,
1061 instrument_id: InstrumentId,
1062 client_order_id: ClientOrderId,
1063 venue_order_id: Option<VenueOrderId>,
1064 reason: &str,
1065 ) {
1066 let ts_event = self.clock.get_time_ns();
1067 self.emitter.emit_order_modify_rejected_event(
1068 strategy_id,
1069 instrument_id,
1070 client_order_id,
1071 venue_order_id,
1072 reason,
1073 ts_event,
1074 );
1075 }
1076
1077 async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
1087 let account_id = self.core.account_id;
1088
1089 if self.core.cache().account(&account_id).is_some() {
1090 log::info!("Account {account_id} registered");
1091 return Ok(());
1092 }
1093
1094 let start = Instant::now();
1095 let timeout = Duration::from_secs_f64(timeout_secs);
1096 let interval = Duration::from_millis(10);
1097
1098 loop {
1099 tokio::time::sleep(interval).await;
1100
1101 if self.core.cache().account(&account_id).is_some() {
1102 log::info!("Account {account_id} registered");
1103 return Ok(());
1104 }
1105
1106 if start.elapsed() >= timeout {
1107 anyhow::bail!(
1108 "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
1109 );
1110 }
1111 }
1112 }
1113}
1114
1115async fn broadcast_partitioned_cancels(
1123 orders: Vec<DydxCancelOrderRequest>,
1124 block_height: u32,
1125 tx_manager: Arc<TransactionManager>,
1126 broadcaster: Arc<TxBroadcaster>,
1127 order_builder: Arc<OrderMessageBuilder>,
1128 emitter: ExecutionEventEmitter,
1129 clock: &'static AtomicTime,
1130) -> anyhow::Result<()> {
1131 if orders.is_empty() {
1132 return Ok(());
1133 }
1134
1135 let (short_term_orders, long_term_orders): (Vec<_>, Vec<_>) = orders
1136 .into_iter()
1137 .partition(|order| order.order_flags == types::ORDER_FLAG_SHORT_TERM);
1138
1139 let mut errors = Vec::new();
1140
1141 if !short_term_orders.is_empty() {
1143 let st_pairs: Vec<_> = short_term_orders
1144 .iter()
1145 .map(|order| (order.instrument_id, order.client_id))
1146 .collect();
1147
1148 log::debug!(
1149 "Batch cancelling {} short-term orders with MsgBatchCancel",
1150 st_pairs.len()
1151 );
1152
1153 match order_builder.build_batch_cancel_short_term(&st_pairs, block_height) {
1154 Ok(msg) => {
1155 let operation = format!("BatchCancel {} short-term orders", st_pairs.len());
1156 match broadcaster
1157 .broadcast_short_term(&tx_manager, vec![msg], &operation)
1158 .await
1159 {
1160 Ok(tx_hash) => {
1161 log::debug!(
1162 "Successfully batch cancelled {} short-term orders, tx_hash: {}",
1163 st_pairs.len(),
1164 tx_hash
1165 );
1166 }
1167 Err(e) => {
1168 let msg = format!("Short-term batch cancel failed: {e:?}");
1169 log::error!("{msg}");
1170 if e.is_definitive_broadcast_rejection() {
1173 emit_partitioned_cancel_rejections(
1174 &short_term_orders,
1175 &emitter,
1176 clock,
1177 &msg,
1178 );
1179 }
1180 errors.push(msg);
1181 }
1182 }
1183 }
1184 Err(e) => {
1185 let msg = format!("Failed to build MsgBatchCancel: {e:?}");
1186 log::error!("{msg}");
1187 errors.push(msg);
1188 }
1189 }
1190 }
1191
1192 if !long_term_orders.is_empty() {
1194 let lt_tuples: Vec<_> = long_term_orders
1195 .iter()
1196 .map(|order| (order.instrument_id, order.client_id, order.order_flags))
1197 .collect();
1198
1199 log::debug!(
1200 "Batch cancelling {} long-term orders",
1201 long_term_orders.len(),
1202 );
1203
1204 match order_builder.build_cancel_orders_batch_with_flags(<_tuples, block_height) {
1205 Ok(cancel_msgs) => {
1206 let operation = format!("BatchCancel {} long-term orders", long_term_orders.len());
1207 match broadcaster
1208 .broadcast_with_retry(&tx_manager, cancel_msgs, &operation)
1209 .await
1210 {
1211 Ok(tx_hash) => {
1212 log::debug!(
1213 "Successfully batch cancelled {} long-term orders, tx_hash: {}",
1214 long_term_orders.len(),
1215 tx_hash
1216 );
1217 }
1218 Err(e) => {
1219 let msg = format!("Long-term batch cancel failed: {e:?}");
1220 log::error!("{msg}");
1221
1222 if e.is_definitive_broadcast_rejection() {
1223 emit_partitioned_cancel_rejections(
1224 &long_term_orders,
1225 &emitter,
1226 clock,
1227 &msg,
1228 );
1229 }
1230 errors.push(msg);
1231 }
1232 }
1233 }
1234 Err(e) => {
1235 let msg = format!("Failed to build long-term cancel messages: {e:?}");
1236 log::error!("{msg}");
1237 errors.push(msg);
1238 }
1239 }
1240 }
1241
1242 if !errors.is_empty() {
1243 anyhow::bail!("partitioned cancel failed: {}", errors.join("; "));
1244 }
1245
1246 Ok(())
1247}
1248
1249fn emit_partitioned_cancel_rejections(
1250 orders: &[DydxCancelOrderRequest],
1251 emitter: &ExecutionEventEmitter,
1252 clock: &AtomicTime,
1253 reason: &str,
1254) {
1255 for order in orders {
1256 emitter.emit_order_cancel_rejected_event(
1257 order.strategy_id,
1258 order.instrument_id,
1259 order.client_order_id,
1260 order.venue_order_id,
1261 reason,
1262 clock.get_time_ns(),
1263 );
1264 }
1265}
1266
1267fn is_definitive_broadcast_rejection(err: &anyhow::Error) -> bool {
1268 err.downcast_ref::<DydxError>()
1269 .is_some_and(DydxError::is_definitive_broadcast_rejection)
1270}
1271
1272#[async_trait(?Send)]
1273impl ExecutionClient for DydxExecutionClient {
1274 fn is_connected(&self) -> bool {
1275 self.core.is_connected()
1276 }
1277
1278 fn client_id(&self) -> ClientId {
1279 self.core.client_id
1280 }
1281
1282 fn account_id(&self) -> AccountId {
1283 self.core.account_id
1284 }
1285
1286 fn venue(&self) -> Venue {
1287 *DYDX_VENUE
1288 }
1289
1290 fn oms_type(&self) -> OmsType {
1291 self.core.oms_type
1292 }
1293
1294 fn get_account(&self) -> Option<AccountAny> {
1295 self.core.cache().account_owned(&self.core.account_id)
1296 }
1297
1298 fn generate_account_state(
1299 &self,
1300 balances: Vec<AccountBalance>,
1301 margins: Vec<MarginBalance>,
1302 reported: bool,
1303 ts_event: UnixNanos,
1304 info: Option<Params>,
1305 ) -> anyhow::Result<()> {
1306 self.emitter
1307 .emit_account_state(balances, margins, reported, ts_event, info);
1308 Ok(())
1309 }
1310
1311 fn start(&mut self) -> anyhow::Result<()> {
1312 if self.core.is_started() {
1313 log::warn!("dYdX execution client already started");
1314 return Ok(());
1315 }
1316
1317 let sender = get_exec_event_sender();
1318 self.emitter.set_sender(sender);
1319 log::info!("Starting dYdX execution client");
1320 self.core.set_started();
1321 Ok(())
1322 }
1323
1324 fn stop(&mut self) -> anyhow::Result<()> {
1325 if self.core.is_stopped() {
1326 log::warn!("dYdX execution client not started");
1327 return Ok(());
1328 }
1329
1330 log::info!("Stopping dYdX execution client");
1331 self.session_tasks.begin_shutdown();
1332 self.begin_pending_shutdown();
1333 self.ws_client.begin_shutdown();
1334 self.core.set_stopped();
1335 self.core.set_disconnected();
1336 Ok(())
1337 }
1338
1339 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
1356 if !self.is_connected() {
1358 let reason = "Cannot submit order: execution client not connected";
1359 log::error!("{reason}");
1360 anyhow::bail!(reason);
1361 }
1362
1363 let current_block = self.block_time_monitor.current_block_height();
1365 let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
1366
1367 let client_order_id = order.client_order_id();
1368 let instrument_id = order.instrument_id();
1369 let strategy_id = order.strategy_id();
1370
1371 if current_block == 0 {
1372 let reason = "Block height not initialized";
1373 log::warn!("Cannot submit order {client_order_id}: {reason}");
1374 self.emitter.emit_order_denied(&order, reason);
1375 return Ok(());
1376 }
1377
1378 if order.is_closed() {
1380 log::warn!("Cannot submit closed order {client_order_id}");
1381 return Ok(());
1382 }
1383
1384 if order.is_quote_quantity() {
1385 let reason = "Quote quantity orders are not supported by dYdX";
1386 log::error!("{reason}");
1387 self.emitter.emit_order_denied(&order, reason);
1388 return Ok(());
1389 }
1390
1391 let unsupported_tif_reason = match order.time_in_force() {
1395 TimeInForce::Fok => Some(
1396 "Fill-or-kill (FOK) orders are deprecated by dYdX v4 (chain rejects with code=48)",
1397 ),
1398 TimeInForce::Day => Some("DAY time-in-force is not supported by dYdX v4"),
1399 _ => None,
1400 };
1401
1402 if let Some(reason) = unsupported_tif_reason {
1403 log::error!("{reason}");
1404 self.emitter.emit_order_denied(&order, reason);
1405 return Ok(());
1406 }
1407
1408 match order.order_type() {
1410 OrderType::Market
1411 | OrderType::Limit
1412 | OrderType::StopMarket
1413 | OrderType::StopLimit
1414 | OrderType::MarketIfTouched
1415 | OrderType::LimitIfTouched => {}
1416 OrderType::TrailingStopMarket | OrderType::TrailingStopLimit => {
1418 let reason = "Trailing stop orders not supported by dYdX v4 protocol";
1419 log::error!("{reason}");
1420 self.emitter.emit_order_denied(&order, reason);
1421 return Ok(());
1422 }
1423 order_type => {
1424 let reason = format!("Order type {order_type:?} not supported by dYdX");
1425 log::error!("{reason}");
1426 self.emitter.emit_order_denied(&order, &reason);
1427 return Ok(());
1428 }
1429 }
1430
1431 let (tx_manager, broadcaster, order_builder) = match self.get_execution_components() {
1433 Ok(components) => components,
1434 Err(e) => {
1435 log::error!("Failed to get execution components: {e}");
1436 self.emitter.emit_order_denied(&order, &e.to_string());
1437 return Ok(());
1438 }
1439 };
1440
1441 let block_height = self.block_time_monitor.current_block_height() as u32;
1442
1443 let encoded = match self.encoder.encode(client_order_id) {
1445 Ok(enc) => enc,
1446 Err(e) => {
1447 log::error!("Failed to generate client order ID: {e}");
1448 self.emitter.emit_order_denied(&order, &e.to_string());
1449 return Ok(());
1450 }
1451 };
1452
1453 self.emitter.emit_order_submitted(&order);
1454 let client_id_u32 = encoded.client_id;
1455 let client_metadata = encoded.client_metadata;
1456
1457 log::debug!(
1458 "[SUBMIT_ORDER] Nautilus '{}' -> dYdX u32={} meta={:#x} | instrument={} side={:?} qty={} type={:?}",
1459 client_order_id,
1460 client_id_u32,
1461 client_metadata,
1462 instrument_id,
1463 order.order_side(),
1464 order.quantity(),
1465 order.order_type()
1466 );
1467
1468 let expire_time = order.expire_time().map(nanos_to_secs_i64);
1470
1471 let order_flags = match order.order_type() {
1473 OrderType::StopMarket
1475 | OrderType::StopLimit
1476 | OrderType::MarketIfTouched
1477 | OrderType::LimitIfTouched => types::ORDER_FLAG_CONDITIONAL,
1478 OrderType::Market => types::ORDER_FLAG_SHORT_TERM,
1480 OrderType::Limit => {
1482 let lifetime = types::OrderLifetime::from_time_in_force(
1483 order.time_in_force(),
1484 expire_time,
1485 false,
1486 order_builder.max_short_term_secs(),
1487 );
1488 lifetime.order_flags()
1489 }
1490 _ => types::ORDER_FLAG_LONG_TERM,
1492 };
1493
1494 let ts_submitted = self.clock.get_time_ns();
1496 let trader_id = order.trader_id();
1497 self.register_order_context(
1498 client_id_u32,
1499 OrderContext {
1500 client_order_id,
1501 trader_id,
1502 strategy_id,
1503 instrument_id,
1504 submitted_at: ts_submitted,
1505 order_flags,
1506 },
1507 );
1508
1509 self.dispatch_state.order_identities.insert(
1512 client_order_id,
1513 OrderIdentity {
1514 instrument_id,
1515 strategy_id,
1516 order_side: order.order_side(),
1517 order_type: order.order_type(),
1518 },
1519 );
1520
1521 self.spawn_order_task(
1522 "submit_order",
1523 strategy_id,
1524 instrument_id,
1525 client_order_id,
1526 async move {
1527 let (msg, order_type_str) = match order.order_type() {
1529 OrderType::Market => {
1530 let msg = order_builder.build_market_order(
1531 instrument_id,
1532 client_id_u32,
1533 client_metadata,
1534 order.order_side(),
1535 order.quantity(),
1536 block_height,
1537 )?;
1538 (msg, "market")
1539 }
1540 OrderType::Limit => {
1541 let msg = order_builder.build_limit_order(
1543 instrument_id,
1544 client_id_u32,
1545 client_metadata,
1546 order.order_side(),
1547 order
1548 .price()
1549 .ok_or_else(|| anyhow::anyhow!("Limit order missing price"))?,
1550 order.quantity(),
1551 order.time_in_force(),
1552 order.is_post_only(),
1553 order.is_reduce_only(),
1554 block_height,
1555 expire_time, )?;
1557 (msg, "limit")
1558 }
1559 OrderType::StopMarket => {
1562 let trigger_price = order.trigger_price().ok_or_else(|| {
1563 anyhow::anyhow!("Stop market order missing trigger_price")
1564 })?;
1565 let cond_expire = order.expire_time().map(nanos_to_secs_i64);
1566 let msg = order_builder.build_stop_market_order(
1567 instrument_id,
1568 client_id_u32,
1569 client_metadata,
1570 order.order_side(),
1571 trigger_price,
1572 order.quantity(),
1573 order.is_reduce_only(),
1574 cond_expire,
1575 )?;
1576 (msg, "stop_market")
1577 }
1578 OrderType::StopLimit => {
1579 let trigger_price = order.trigger_price().ok_or_else(|| {
1580 anyhow::anyhow!("Stop limit order missing trigger_price")
1581 })?;
1582 let limit_price = order.price().ok_or_else(|| {
1583 anyhow::anyhow!("Stop limit order missing limit price")
1584 })?;
1585 let cond_expire = order.expire_time().map(nanos_to_secs_i64);
1586 let msg = order_builder.build_stop_limit_order(
1587 instrument_id,
1588 client_id_u32,
1589 client_metadata,
1590 order.order_side(),
1591 trigger_price,
1592 limit_price,
1593 order.quantity(),
1594 order.time_in_force(),
1595 order.is_post_only(),
1596 order.is_reduce_only(),
1597 cond_expire,
1598 )?;
1599 (msg, "stop_limit")
1600 }
1601 OrderType::MarketIfTouched => {
1603 let trigger_price = order.trigger_price().ok_or_else(|| {
1604 anyhow::anyhow!("Take profit market order missing trigger_price")
1605 })?;
1606 let cond_expire = order.expire_time().map(nanos_to_secs_i64);
1607 let msg = order_builder.build_take_profit_market_order(
1608 instrument_id,
1609 client_id_u32,
1610 client_metadata,
1611 order.order_side(),
1612 trigger_price,
1613 order.quantity(),
1614 order.is_reduce_only(),
1615 cond_expire,
1616 )?;
1617 (msg, "take_profit_market")
1618 }
1619 OrderType::LimitIfTouched => {
1621 let trigger_price = order.trigger_price().ok_or_else(|| {
1622 anyhow::anyhow!("Take profit limit order missing trigger_price")
1623 })?;
1624 let limit_price = order.price().ok_or_else(|| {
1625 anyhow::anyhow!("Take profit limit order missing limit price")
1626 })?;
1627 let cond_expire = order.expire_time().map(nanos_to_secs_i64);
1628 let msg = order_builder.build_take_profit_limit_order(
1629 instrument_id,
1630 client_id_u32,
1631 client_metadata,
1632 order.order_side(),
1633 trigger_price,
1634 limit_price,
1635 order.quantity(),
1636 order.time_in_force(),
1637 order.is_post_only(),
1638 order.is_reduce_only(),
1639 cond_expire,
1640 )?;
1641 (msg, "take_profit_limit")
1642 }
1643 _ => unreachable!("Order type already validated"),
1644 };
1645
1646 let operation = format!("Submit {order_type_str} order {client_order_id}");
1649
1650 if order_flags == types::ORDER_FLAG_SHORT_TERM {
1651 broadcaster
1652 .broadcast_short_term(&tx_manager, vec![msg], &operation)
1653 .await?;
1654 } else {
1655 broadcaster
1656 .broadcast_with_retry(&tx_manager, vec![msg], &operation)
1657 .await?;
1658 }
1659 log::debug!("Successfully submitted {order_type_str} order: {client_order_id}");
1660
1661 Ok(())
1662 },
1663 );
1664
1665 Ok(())
1666 }
1667
1668 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1669 let orders = self.core.get_orders_for_list(&cmd.order_list)?;
1670 let order_count = orders.len();
1671
1672 if !self.is_connected() {
1674 let reason = "Cannot submit order list: execution client not connected";
1675 log::error!("{reason}");
1676 anyhow::bail!(reason);
1677 }
1678
1679 let unsupported_tif_reason = |order: &OrderAny| match order.time_in_force() {
1684 TimeInForce::Fok => Some(
1685 "Fill-or-kill (FOK) orders are deprecated by dYdX v4 (chain rejects with code=48)",
1686 ),
1687 TimeInForce::Day => Some("DAY time-in-force is not supported by dYdX v4"),
1688 _ => None,
1689 };
1690
1691 if orders
1692 .iter()
1693 .any(|order| unsupported_tif_reason(order).is_some())
1694 {
1695 for order in &orders {
1698 let reason = unsupported_tif_reason(order)
1699 .unwrap_or("Order list denied: sibling order has unsupported time in force");
1700 self.emitter.emit_order_denied(order, reason);
1701 }
1702 return Ok(());
1703 }
1704
1705 let current_block = self.block_time_monitor.current_block_height();
1707 if current_block == 0 {
1708 let reason = "Block height not initialized";
1709 log::warn!("Cannot submit order list: {reason}");
1710 for order in &orders {
1711 self.emitter.emit_order_denied(order, reason);
1712 }
1713 return Ok(());
1714 }
1715
1716 let (tx_manager, broadcaster, order_builder) = match self.get_execution_components() {
1718 Ok(components) => components,
1719 Err(e) => {
1720 log::error!("Failed to get execution components for batch: {e}");
1721 for order in &orders {
1722 self.emitter.emit_order_denied(order, &e.to_string());
1723 }
1724 return Ok(());
1725 }
1726 };
1727
1728 let mut order_params: Vec<LimitOrderParams> = Vec::with_capacity(order_count);
1730 let mut order_info: Vec<(ClientOrderId, InstrumentId, StrategyId)> =
1731 Vec::with_capacity(order_count);
1732
1733 for order in &orders {
1734 if order.order_type() != OrderType::Limit {
1736 log::warn!(
1737 "Order {} has type {:?}, falling back to individual submission",
1738 order.client_order_id(),
1739 order.order_type()
1740 );
1741 let submit_cmd = SubmitOrder::new(
1743 cmd.trader_id,
1744 cmd.client_id,
1745 cmd.strategy_id,
1746 order.instrument_id(),
1747 order.client_order_id(),
1748 order.init_event().clone(),
1749 cmd.exec_algorithm_id,
1750 cmd.position_id,
1751 cmd.params.clone(),
1752 UUID4::new(),
1753 cmd.ts_init,
1754 cmd.correlation_id,
1755 );
1756
1757 if let Err(e) = self.submit_order(submit_cmd) {
1758 log::error!(
1759 "Failed to submit order {} from order list: {e}",
1760 order.client_order_id()
1761 );
1762 }
1763 continue;
1764 }
1765
1766 if order.is_quote_quantity() {
1767 self.emitter
1768 .emit_order_denied(order, "Quote quantity orders are not supported by dYdX");
1769 continue;
1770 }
1771
1772 let Some(price) = order.price() else {
1774 self.emitter
1775 .emit_order_denied(order, "Limit order missing price");
1776 continue;
1777 };
1778
1779 let encoded = match self.encoder.encode(order.client_order_id()) {
1781 Ok(enc) => enc,
1782 Err(e) => {
1783 log::error!("Failed to generate client order ID: {e}");
1784 self.emitter.emit_order_denied(order, &e.to_string());
1785 continue;
1786 }
1787 };
1788 let client_id_u32 = encoded.client_id;
1789 let client_metadata = encoded.client_metadata;
1790
1791 self.emitter.emit_order_submitted(order);
1793
1794 let expire_time_secs = order.expire_time().map(nanos_to_secs_i64);
1796 let lifetime = types::OrderLifetime::from_time_in_force(
1797 order.time_in_force(),
1798 expire_time_secs,
1799 false,
1800 order_builder.max_short_term_secs(),
1801 );
1802
1803 let ts_submitted = self.clock.get_time_ns();
1805 self.register_order_context(
1806 client_id_u32,
1807 OrderContext {
1808 client_order_id: order.client_order_id(),
1809 trader_id: order.trader_id(),
1810 strategy_id: order.strategy_id(),
1811 instrument_id: order.instrument_id(),
1812 submitted_at: ts_submitted,
1813 order_flags: lifetime.order_flags(),
1814 },
1815 );
1816
1817 self.dispatch_state.order_identities.insert(
1819 order.client_order_id(),
1820 OrderIdentity {
1821 instrument_id: order.instrument_id(),
1822 strategy_id: order.strategy_id(),
1823 order_side: order.order_side(),
1824 order_type: order.order_type(),
1825 },
1826 );
1827
1828 order_params.push(LimitOrderParams {
1830 instrument_id: order.instrument_id(),
1831 client_order_id: client_id_u32,
1832 client_metadata,
1833 side: order.order_side(),
1834 price,
1835 quantity: order.quantity(),
1836 time_in_force: order.time_in_force(),
1837 post_only: order.is_post_only(),
1838 reduce_only: order.is_reduce_only(),
1839 expire_time_ns: order.expire_time(),
1840 });
1841 order_info.push((
1842 order.client_order_id(),
1843 order.instrument_id(),
1844 order.strategy_id(),
1845 ));
1846 }
1847
1848 if order_params.is_empty() {
1850 return Ok(());
1851 }
1852
1853 let has_short_term = order_params
1857 .iter()
1858 .any(|params| order_builder.is_short_term_order(params));
1859
1860 let block_height = current_block as u32;
1861 let emitter = self.emitter.clone();
1862 let clock = self.clock;
1863
1864 if has_short_term {
1865 log::debug!(
1867 "Submitting {} short-term limit orders concurrently (sequence not consumed)",
1868 order_params.len()
1869 );
1870
1871 self.spawn_labeled("batch_submit_short_term", async move {
1872 let mut handles = tokio::task::JoinSet::new();
1875
1876 for (params, (client_order_id, instrument_id, strategy_id)) in
1877 order_params.into_iter().zip(order_info)
1878 {
1879 let tx_manager = tx_manager.clone();
1880 let broadcaster = broadcaster.clone();
1881 let order_builder = order_builder.clone();
1882 let emitter = emitter.clone();
1883
1884 handles.spawn(async move {
1885 let msg = match order_builder
1887 .build_limit_order_from_params(¶ms, block_height)
1888 {
1889 Ok(m) => m,
1890 Err(e) => {
1891 log::error!(
1894 "Failed to build order message for {client_order_id}: {e:?}"
1895 );
1896 return;
1897 }
1898 };
1899
1900 let operation = format!("Submit short-term order {client_order_id}");
1902
1903 if let Err(e) = broadcaster
1904 .broadcast_short_term(&tx_manager, vec![msg], &operation)
1905 .await
1906 {
1907 if e.is_definitive_broadcast_rejection() {
1908 let error_msg = format!("Order submission failed: {e:?}");
1909 log::error!("{error_msg}");
1910 let ts_event = clock.get_time_ns();
1911 emitter.emit_order_rejected_event(
1912 strategy_id,
1913 instrument_id,
1914 client_order_id,
1915 &error_msg,
1916 ts_event,
1917 false,
1918 );
1919 } else {
1920 log::warn!(
1921 "Ambiguous dYdX submit failure for {client_order_id}, awaiting reconciliation: {e:?}"
1922 );
1923 }
1924 }
1925 });
1926 }
1927
1928 while let Some(result) = handles.join_next().await {
1930 if let Err(e) = result
1931 && !e.is_cancelled()
1932 {
1933 log::warn!("dYdX short-term order task failed: {e}");
1934 }
1935 }
1936 });
1937 } else {
1938 log::debug!(
1940 "Batch submitting {} long-term limit orders in single transaction",
1941 order_params.len()
1942 );
1943
1944 self.spawn_labeled("batch_submit_long_term", async move {
1945 let msgs: Result<Vec<_>, _> = order_params
1947 .iter()
1948 .map(|params| order_builder.build_limit_order_from_params(params, block_height))
1949 .collect();
1950
1951 let msgs = match msgs {
1952 Ok(m) => m,
1953 Err(e) => {
1954 log::error!("Failed to build batch order messages: {e:?}");
1957 return;
1958 }
1959 };
1960
1961 let operation = format!("Submit batch of {} limit orders", msgs.len());
1963
1964 if let Err(e) = broadcaster
1965 .broadcast_with_retry(&tx_manager, msgs, &operation)
1966 .await
1967 {
1968 if e.is_definitive_broadcast_rejection() {
1969 let error_msg = format!("Batch order submission failed: {e:?}");
1972 log::error!("{error_msg}");
1973 let ts_event = clock.get_time_ns();
1974
1975 for (client_order_id, instrument_id, strategy_id) in order_info {
1976 emitter.emit_order_rejected_event(
1977 strategy_id,
1978 instrument_id,
1979 client_order_id,
1980 &error_msg,
1981 ts_event,
1982 false,
1983 );
1984 }
1985 } else {
1986 log::warn!(
1987 "Ambiguous dYdX batch submit failure, awaiting reconciliation: {e:?}"
1988 );
1989 }
1990 }
1991 });
1992 }
1993
1994 Ok(())
1995 }
1996
1997 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
2001 let reason = "dYdX does not support order modification. Use cancel and resubmit instead.";
2002 log::error!("{reason}");
2003
2004 self.send_modify_rejected(
2005 cmd.strategy_id,
2006 cmd.instrument_id,
2007 cmd.client_order_id,
2008 cmd.venue_order_id,
2009 reason,
2010 );
2011 Ok(())
2012 }
2013
2014 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
2033 if !self.is_connected() {
2034 anyhow::bail!("Cannot cancel order: not connected");
2035 }
2036
2037 let client_order_id = cmd.client_order_id;
2038 let instrument_id = cmd.instrument_id;
2039 let strategy_id = cmd.strategy_id;
2040 let venue_order_id = cmd.venue_order_id;
2041
2042 let (order_time_in_force, order_expire_time) = {
2043 let cache = self.core.cache();
2044
2045 let order = match cache.order(&client_order_id) {
2046 Some(order) => order,
2047 None => {
2048 log::error!("Cannot cancel order {client_order_id}: not found in cache");
2049 return Ok(()); }
2051 };
2052
2053 if order.is_closed() {
2055 log::warn!(
2056 "CancelOrder command for {} when order already {} (will not send to exchange)",
2057 client_order_id,
2058 order.status()
2059 );
2060 return Ok(());
2061 }
2062
2063 if cache.instrument(&instrument_id).is_none() {
2065 log::error!(
2066 "Cannot cancel order {client_order_id}: instrument {instrument_id} not found in cache"
2067 );
2068 return Ok(()); }
2070
2071 (
2073 order.time_in_force(),
2074 order.expire_time().map(nanos_to_secs_i64),
2075 )
2076 }; log::debug!("Cancelling order {client_order_id} for instrument {instrument_id}");
2079
2080 let (tx_manager, broadcaster, order_builder) = match self.get_execution_components() {
2082 Ok(components) => components,
2083 Err(e) => {
2084 log::error!("Failed to get execution components for cancel: {e}");
2085 return Ok(());
2086 }
2087 };
2088
2089 let block_height = self.block_time_monitor.current_block_height() as u32;
2090
2091 let encoded = match self.encoder.get(&client_order_id) {
2093 Some(enc) => enc,
2094 None => {
2095 log::error!("Client order ID {client_order_id} not found in cache");
2096 anyhow::bail!("Client order ID not found in cache")
2097 }
2098 };
2099 let client_id_u32 = encoded.client_id;
2100
2101 log::debug!(
2102 "[CANCEL_ORDER] Nautilus '{client_order_id}' -> dYdX u32={client_id_u32} | instrument={instrument_id}"
2103 );
2104
2105 let order_flags = self.get_order_context(client_id_u32).map_or_else(
2108 || {
2109 log::warn!(
2111 "Order context not found for {client_order_id}, deriving flags from order"
2112 );
2113 types::OrderLifetime::from_time_in_force(
2114 order_time_in_force, order_expire_time, false,
2117 order_builder.max_short_term_secs(),
2118 )
2119 .order_flags()
2120 },
2121 |ctx| ctx.order_flags,
2122 );
2123
2124 let clock = self.clock;
2125 let emitter = self.emitter.clone();
2126
2127 self.spawn_task("cancel_order", async move {
2128 let cancel_msg = match order_builder.build_cancel_order_with_flags(
2130 instrument_id,
2131 client_id_u32,
2132 order_flags,
2133 block_height,
2134 ) {
2135 Ok(msg) => msg,
2136 Err(e) => {
2137 log::warn!(
2140 "Cancel command failed local validation for {client_order_id}: {e:?}"
2141 );
2142 return Ok(());
2143 }
2144 };
2145
2146 let cancel_op = format!("Cancel order {client_order_id}");
2148 let result = if order_flags == types::ORDER_FLAG_SHORT_TERM {
2149 broadcaster
2150 .broadcast_short_term(&tx_manager, vec![cancel_msg], &cancel_op)
2151 .await
2152 } else {
2153 broadcaster
2154 .broadcast_with_retry(&tx_manager, vec![cancel_msg], &cancel_op)
2155 .await
2156 };
2157
2158 match result {
2159 Ok(_) => {
2160 log::debug!("Successfully cancelled order: {client_order_id}");
2161 }
2162 Err(e) if e.is_definitive_broadcast_rejection() => {
2163 log::error!("Failed to cancel order {client_order_id}: {e:?}");
2164
2165 let ts_event = clock.get_time_ns();
2166 emitter.emit_order_cancel_rejected_event(
2167 strategy_id,
2168 instrument_id,
2169 client_order_id,
2170 venue_order_id,
2171 &format!("Cancel order failed: {e:?}"),
2172 ts_event,
2173 );
2174 }
2175 Err(e) => {
2176 log::warn!(
2177 "Ambiguous dYdX cancel failure for {client_order_id}, awaiting reconciliation: {e:?}"
2178 );
2179 }
2180 }
2181
2182 Ok(())
2183 });
2184
2185 Ok(())
2186 }
2187
2188 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
2189 if !self.is_connected() {
2190 anyhow::bail!("Cannot cancel orders: not connected");
2191 }
2192
2193 let instrument_id = cmd.instrument_id;
2194 let order_side_filter = cmd.order_side;
2195
2196 let order_data: Vec<CancelAllOrderData> = {
2197 let cache = self.core.cache();
2198 let side_filter = order_side_filter;
2199 cache
2200 .orders_open(None, Some(&instrument_id), None, None, side_filter)
2201 .into_iter()
2202 .map(|order| {
2203 (
2204 order.strategy_id(),
2205 order.client_order_id(),
2206 order.venue_order_id(),
2207 order.time_in_force(),
2208 order.expire_time(),
2209 )
2210 })
2211 .collect()
2212 }; let short_term_count = order_data
2216 .iter()
2217 .filter(|(_, _, _, tif, _)| matches!(tif, TimeInForce::Ioc | TimeInForce::Fok))
2218 .count();
2219 let long_term_count = order_data.len() - short_term_count;
2220
2221 log::debug!(
2222 "Cancel all orders: total={}, short_term={}, long_term={}, instrument_id={instrument_id}, order_side={order_side_filter:?}",
2223 order_data.len(),
2224 short_term_count,
2225 long_term_count
2226 );
2227
2228 let (tx_manager, broadcaster, order_builder) = match self.get_execution_components() {
2230 Ok(components) => components,
2231 Err(e) => {
2232 log::error!("Failed to get execution components for cancel_all: {e}");
2233 return Ok(());
2234 }
2235 };
2236
2237 let block_height = self.block_time_monitor.current_block_height() as u32;
2238
2239 let mut orders_to_cancel = Vec::new();
2242
2243 for (strategy_id, client_order_id, venue_order_id, _time_in_force, _expire_time) in
2244 &order_data
2245 {
2246 let Some(encoded) = self.encoder.get(client_order_id) else {
2247 log::warn!("Cannot cancel order {client_order_id}: not found in encoder");
2248 continue;
2249 };
2250 let client_id_u32 = encoded.client_id;
2251
2252 let Some(ctx) = self.get_order_context(client_id_u32) else {
2254 log::debug!(
2255 "Skipping cancel for {client_order_id}: order context already cleaned up (terminal)"
2256 );
2257 continue;
2258 };
2259 orders_to_cancel.push(DydxCancelOrderRequest {
2260 instrument_id,
2261 client_id: client_id_u32,
2262 order_flags: ctx.order_flags,
2263 strategy_id: *strategy_id,
2264 client_order_id: *client_order_id,
2265 venue_order_id: *venue_order_id,
2266 });
2267 }
2268
2269 if orders_to_cancel.is_empty() {
2270 return Ok(());
2271 }
2272
2273 log::debug!(
2274 "Cancel all: {} orders (short_term={}, long_term={}), instrument_id={instrument_id}, order_side={order_side_filter:?}",
2275 orders_to_cancel.len(),
2276 orders_to_cancel
2277 .iter()
2278 .filter(|order| order.order_flags == types::ORDER_FLAG_SHORT_TERM)
2279 .count(),
2280 orders_to_cancel
2281 .iter()
2282 .filter(|order| order.order_flags != types::ORDER_FLAG_SHORT_TERM)
2283 .count(),
2284 );
2285
2286 let clock = self.clock;
2287 let emitter = self.emitter.clone();
2288
2289 self.spawn_task("cancel_all_orders", async move {
2290 broadcast_partitioned_cancels(
2291 orders_to_cancel,
2292 block_height,
2293 tx_manager,
2294 broadcaster,
2295 order_builder,
2296 emitter,
2297 clock,
2298 )
2299 .await
2300 });
2301
2302 Ok(())
2303 }
2304
2305 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
2306 if cmd.cancels.is_empty() {
2307 return Ok(());
2308 }
2309
2310 if !self.is_connected() {
2311 anyhow::bail!("Cannot cancel orders: not connected");
2312 }
2313
2314 let (tx_manager, broadcaster, order_builder) = match self.get_execution_components() {
2316 Ok(components) => components,
2317 Err(e) => {
2318 log::error!("Failed to get execution components for batch cancel: {e}");
2319 return Ok(());
2320 }
2321 };
2322
2323 let mut orders_to_cancel = Vec::with_capacity(cmd.cancels.len());
2325 for cancel in &cmd.cancels {
2326 let client_order_id = cancel.client_order_id;
2327 let encoded = match self.encoder.get(&client_order_id) {
2328 Some(enc) => enc,
2329 None => {
2330 log::warn!(
2331 "No u32 mapping found for client_order_id={client_order_id}, skipping cancel"
2332 );
2333 continue;
2334 }
2335 };
2336 let client_id_u32 = encoded.client_id;
2337
2338 let Some(ctx) = self.get_order_context(client_id_u32) else {
2340 log::debug!(
2341 "Skipping cancel for {client_order_id}: order context already cleaned up (terminal)"
2342 );
2343 continue;
2344 };
2345
2346 orders_to_cancel.push(DydxCancelOrderRequest {
2347 instrument_id: cancel.instrument_id,
2348 client_id: client_id_u32,
2349 order_flags: ctx.order_flags,
2350 strategy_id: cancel.strategy_id,
2351 client_order_id,
2352 venue_order_id: cancel.venue_order_id,
2353 });
2354 }
2355
2356 if orders_to_cancel.is_empty() {
2357 log::warn!("No valid orders to cancel in batch");
2358 return Ok(());
2359 }
2360
2361 let block_height = self.block_time_monitor.current_block_height() as u32;
2362 let clock = self.clock;
2363 let emitter = self.emitter.clone();
2364
2365 log::debug!(
2366 "Batch cancelling {} orders via partitioned strategy",
2367 orders_to_cancel.len(),
2368 );
2369
2370 self.spawn_task("batch_cancel_orders", async move {
2371 broadcast_partitioned_cancels(
2372 orders_to_cancel,
2373 block_height,
2374 tx_manager,
2375 broadcaster,
2376 order_builder,
2377 emitter,
2378 clock,
2379 )
2380 .await
2381 });
2382
2383 Ok(())
2384 }
2385
2386 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
2387 let http_client = self.http_client.clone();
2388 let wallet_address = self.wallet_address.clone();
2389 let subaccount_number = self.subaccount_number;
2390 let account_id = self.core.account_id;
2391 let emitter = self.emitter.clone();
2392
2393 self.spawn_task("query_account", async move {
2394 let account_state = http_client
2395 .request_account_state(&wallet_address, subaccount_number, account_id)
2396 .await
2397 .context("failed to query account state")?;
2398
2399 emitter.emit_account_state(
2400 account_state.balances.clone(),
2401 account_state.margins.clone(),
2402 account_state.is_reported,
2403 account_state.ts_event,
2404 account_state.info,
2405 );
2406 Ok(())
2407 });
2408
2409 Ok(())
2410 }
2411
2412 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
2413 log::debug!("Querying order: client_order_id={}", cmd.client_order_id);
2414
2415 let http_client = self.http_client.clone();
2416 let wallet_address = self.wallet_address.clone();
2417 let subaccount_number = self.subaccount_number;
2418 let account_id = self.core.account_id;
2419 let emitter = self.emitter.clone();
2420 let client_order_id = cmd.client_order_id;
2421 let venue_order_id = cmd.venue_order_id;
2422 let instrument_id = cmd.instrument_id;
2423
2424 self.spawn_task("query_order", async move {
2425 let reports = http_client
2426 .request_order_status_reports(
2427 &wallet_address,
2428 subaccount_number,
2429 account_id,
2430 Some(instrument_id),
2431 )
2432 .await
2433 .context("failed to query order status")?;
2434
2435 let report = reports.into_iter().find(|r| {
2437 if venue_order_id.is_some_and(|vid| r.venue_order_id == vid) {
2438 return true;
2439 }
2440 r.client_order_id.is_some_and(|cid| cid == client_order_id)
2441 });
2442
2443 if let Some(report) = report {
2444 emitter.send_order_status_report(report);
2445 } else {
2446 log::warn!(
2447 "No order found for client_order_id={client_order_id}, venue_order_id={venue_order_id:?}"
2448 );
2449 }
2450
2451 Ok(())
2452 });
2453
2454 Ok(())
2455 }
2456
2457 async fn connect(&mut self) -> anyhow::Result<()> {
2458 if self.core.is_connected() && self.session_tasks.is_open() && self.pending_tasks.is_open()
2459 {
2460 log::warn!("dYdX execution client already connected");
2461 return Ok(());
2462 }
2463
2464 log::info!("Connecting to dYdX");
2465
2466 self.prepare_task_groups().await?;
2467 let ws_client = self.ws_client.clone();
2468 let setup_guard =
2469 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
2470 ws_client.begin_shutdown();
2471 });
2472
2473 log::debug!("Loading instruments from HTTP API");
2474 self.http_client.fetch_and_cache_instruments().await?;
2475 log::debug!(
2476 "Loaded {} instruments from HTTP into shared cache",
2477 self.http_client.cached_instruments_count()
2478 );
2479 self.mark_instruments_initialized();
2480
2481 let grpc_urls = self.config.get_grpc_urls();
2483 let mut grpc_client = DydxGrpcClient::new_with_fallback(&grpc_urls)
2484 .await
2485 .context("failed to construct dYdX gRPC client")?;
2486 log::debug!("gRPC client initialized");
2487
2488 let initial_height = grpc_client
2490 .latest_block_height()
2491 .await
2492 .context("failed to fetch initial block height")?;
2493 self.block_time_monitor
2495 .record_block(initial_height.0 as u64, jiff::Timestamp::now());
2496 log::debug!("Initial block height: {}", initial_height.0);
2497
2498 *self.grpc_client.write().await = Some(grpc_client.clone());
2499
2500 let private_key =
2502 Self::resolve_private_key(&self.config).context("failed to resolve private key")?;
2503 let tx_manager = Arc::new(
2504 TransactionManager::new(
2505 grpc_client.clone(),
2506 &private_key,
2507 self.wallet_address.clone(),
2508 self.get_chain_id(),
2509 )
2510 .context("failed to create TransactionManager")?,
2511 );
2512
2513 tx_manager
2514 .resolve_authenticators()
2515 .await
2516 .context("failed to resolve authenticators")?;
2517
2518 tx_manager
2521 .initialize_sequence()
2522 .await
2523 .context("failed to initialize sequence")?;
2524
2525 self.tx_manager = Some(tx_manager);
2526 self.broadcaster = Some(Arc::new(TxBroadcaster::new(
2527 grpc_client,
2528 self.config.grpc_quota(),
2529 )));
2530 self.order_builder = Some(Arc::new(OrderMessageBuilder::new(
2531 self.http_client.clone(),
2532 self.wallet_address.clone(),
2533 self.subaccount_number,
2534 self.block_time_monitor.clone(),
2535 )));
2536 log::debug!(
2537 "OrderMessageBuilder initialized (block_time_monitor ready: {}, max_short_term: {:.1}s)",
2538 self.block_time_monitor.is_ready(),
2539 SHORT_TERM_ORDER_MAXIMUM_LIFETIME as f64
2540 * self.block_time_monitor.seconds_per_block_or_default()
2541 );
2542
2543 let session_result = async {
2545 self.ws_client.connect().await?;
2546 log::debug!("WebSocket connected");
2547
2548 self.ws_client.subscribe_block_height().await?;
2549 log::debug!("Subscribed to block height updates");
2550
2551 self.ws_client.subscribe_markets().await?;
2552 log::debug!("Subscribed to markets");
2553
2554 log::debug!(
2555 "Using wallet address for queries: {} (subaccount {})",
2556 self.wallet_address,
2557 self.subaccount_number
2558 );
2559 self.ws_client
2560 .subscribe_subaccount(&self.wallet_address, self.subaccount_number)
2561 .await?;
2562 log::debug!(
2563 "Subscribed to subaccount updates: {}/{}",
2564 self.wallet_address,
2565 self.subaccount_number
2566 );
2567
2568 let stream = self.ws_client.stream();
2569 self.spawn_ws_stream_handler(stream)?;
2570
2571 self.await_account_registered(30.0).await?;
2575
2576 Ok::<(), anyhow::Error>(())
2577 }
2578 .await;
2579
2580 if let Err(e) = session_result {
2581 if let Err(teardown_error) = self.teardown_partial_connect().await {
2582 return Err(e.context(format!(
2583 "dYdX execution startup teardown failed: {teardown_error}"
2584 )));
2585 }
2586 return Err(e);
2587 }
2588
2589 self.core.set_connected();
2590 setup_guard.disarm();
2591 log::info!("Connected: client_id={}", self.core.client_id);
2592 Ok(())
2593 }
2594
2595 async fn disconnect(&mut self) -> anyhow::Result<()> {
2596 log::info!("Disconnecting from dYdX");
2597
2598 self.begin_pending_shutdown();
2599 let ws_client = self.ws_client.clone();
2600 let disconnect_guard = TaskGroupGuard::new(&[&self.session_tasks], move || {
2601 ws_client.begin_shutdown();
2602 });
2603
2604 let _ = self
2606 .ws_client
2607 .unsubscribe_subaccount(&self.wallet_address, self.subaccount_number)
2608 .await
2609 .map_err(|e| log::warn!("Failed to unsubscribe from subaccount: {e}"));
2610
2611 let _ = self
2613 .ws_client
2614 .unsubscribe_markets()
2615 .await
2616 .map_err(|e| log::warn!("Failed to unsubscribe from markets: {e}"));
2617
2618 let _ = self
2620 .ws_client
2621 .unsubscribe_block_height()
2622 .await
2623 .map_err(|e| log::warn!("Failed to unsubscribe from block height: {e}"));
2624
2625 let shutdown_result = self.teardown_partial_connect().await;
2627 disconnect_guard.disarm();
2628 shutdown_result?;
2629 log::info!("Disconnected: client_id={}", self.core.client_id);
2630 Ok(())
2631 }
2632
2633 async fn generate_order_status_report(
2634 &self,
2635 cmd: &GenerateOrderStatusReport,
2636 ) -> anyhow::Result<Option<OrderStatusReport>> {
2637 let market = cmd
2642 .instrument_id
2643 .map(|id| id.symbol.as_str().trim_end_matches("-PERP").to_string());
2644
2645 let response = self
2646 .http_client
2647 .inner
2648 .get_orders(
2649 &self.wallet_address,
2650 self.subaccount_number,
2651 market.as_deref(),
2652 Some(DYDX_INDEXER_REPORT_LIMIT),
2653 )
2654 .await
2655 .context("failed to fetch order from dYdX API")?;
2656
2657 if response.is_empty() {
2658 log::debug!(
2659 "No orders returned for {}/subaccount={} (market_filter={:?})",
2660 self.wallet_address,
2661 self.subaccount_number,
2662 market,
2663 );
2664 return Ok(None);
2665 }
2666
2667 let ts_init = UnixNanos::default();
2668 let scanned_count = response.len();
2669
2670 let report = find_matching_order_report(
2671 &response,
2672 cmd.instrument_id,
2673 cmd.client_order_id,
2674 cmd.venue_order_id,
2675 |clob_pair_id| self.get_instrument_by_clob_pair_id(clob_pair_id),
2676 &self.encoder,
2677 self.core.account_id,
2678 ts_init,
2679 )?;
2680
2681 if report.is_none() {
2682 let page_full = scanned_count == DYDX_INDEXER_REPORT_LIMIT as usize;
2686 log::debug!(
2687 "No order matched filters for {}/subaccount={} \
2688 (client_order_id={:?}, venue_order_id={:?}, instrument_id={:?}, \
2689 scanned={scanned_count}, page_full={page_full}, limit={DYDX_INDEXER_REPORT_LIMIT})",
2690 self.wallet_address,
2691 self.subaccount_number,
2692 cmd.client_order_id,
2693 cmd.venue_order_id,
2694 cmd.instrument_id,
2695 );
2696 }
2697
2698 Ok(report)
2699 }
2700
2701 async fn generate_order_status_reports(
2702 &self,
2703 cmd: &GenerateOrderStatusReports,
2704 ) -> anyhow::Result<Vec<OrderStatusReport>> {
2705 let response = self
2707 .http_client
2708 .inner
2709 .get_orders(
2710 &self.wallet_address,
2711 self.subaccount_number,
2712 None, Some(DYDX_INDEXER_REPORT_LIMIT),
2714 )
2715 .await
2716 .context("failed to fetch orders from dYdX API")?;
2717
2718 let mut reports = Vec::new();
2719 let ts_init = UnixNanos::default();
2720
2721 for order in response {
2722 let instrument = match self.get_instrument_by_clob_pair_id(order.clob_pair_id) {
2723 Some(inst) => inst,
2724 None => continue,
2725 };
2726
2727 if let Some(filter_id) = cmd.instrument_id
2728 && instrument.id() != filter_id
2729 {
2730 continue;
2731 }
2732
2733 match parse_order_status_report(&order, &instrument, self.core.account_id, ts_init) {
2734 Ok(mut r) => {
2735 if !order.client_id.is_empty()
2736 && let Ok(client_id_u32) = order.client_id.parse::<u32>()
2737 {
2738 self.encoder.register_known_client_id(client_id_u32);
2739
2740 if let Some(decoded) = self
2741 .encoder
2742 .decode_if_known(client_id_u32, order.client_metadata)
2743 {
2744 log::debug!(
2745 "Decoded order: dYdX client_id={} meta={:#x} -> '{}'",
2746 client_id_u32,
2747 order.client_metadata,
2748 decoded,
2749 );
2750 r.client_order_id = Some(decoded);
2751 }
2752 }
2753 reports.push(r);
2754 }
2755 Err(e) => {
2756 log::warn!("Failed to parse order status report: {e}");
2757 }
2758 }
2759 }
2760
2761 if cmd.open_only {
2763 reports.retain(|r| r.order_status.is_open());
2764 }
2765
2766 if let Some(start) = cmd.start {
2768 reports.retain(|r| r.ts_last >= start);
2769 }
2770
2771 if let Some(end) = cmd.end {
2772 reports.retain(|r| r.ts_last <= end);
2773 }
2774
2775 {
2782 let cache = self.core.cache();
2783 reports.retain(|r| {
2784 let Some(cid) = r.client_order_id else {
2785 return true;
2786 };
2787
2788 match cache.order(&cid) {
2789 Some(order) if order.status().is_closed() => {
2790 log::debug!(
2791 "Skipping reconciliation report for terminal order {cid} \
2792 (cache status={:?}, venue status={:?})",
2793 order.status(),
2794 r.order_status,
2795 );
2796 false
2797 }
2798 _ => true,
2799 }
2800 });
2801 }
2802
2803 Ok(reports)
2804 }
2805
2806 async fn generate_fill_reports(
2807 &self,
2808 cmd: GenerateFillReports,
2809 ) -> anyhow::Result<Vec<FillReport>> {
2810 let response = self
2811 .http_client
2812 .inner
2813 .get_fills(
2814 &self.wallet_address,
2815 self.subaccount_number,
2816 None, Some(DYDX_INDEXER_REPORT_LIMIT),
2818 )
2819 .await
2820 .context("failed to fetch fills from dYdX API")?;
2821
2822 let mut reports = Vec::new();
2823 let ts_init = UnixNanos::default();
2824
2825 for fill in response.fills {
2826 let instrument = match self.get_instrument_by_market(&fill.market) {
2827 Some(inst) => inst,
2828 None => {
2829 log::warn!("Unknown market in fill: {}", fill.market);
2830 continue;
2831 }
2832 };
2833
2834 if let Some(filter_id) = cmd.instrument_id
2835 && instrument.id() != filter_id
2836 {
2837 continue;
2838 }
2839
2840 let report = match parse_fill_report(&fill, &instrument, self.core.account_id, ts_init)
2841 {
2842 Ok(r) => r,
2843 Err(e) => {
2844 log::warn!("Failed to parse fill report: {e}");
2845 continue;
2846 }
2847 };
2848
2849 reports.push(report);
2850 }
2851
2852 if let Some(venue_order_id) = cmd.venue_order_id {
2853 reports.retain(|r| r.venue_order_id.as_str() == venue_order_id.as_str());
2854 }
2855
2856 Ok(reports)
2857 }
2858
2859 async fn generate_position_status_reports(
2860 &self,
2861 cmd: &GeneratePositionStatusReports,
2862 ) -> anyhow::Result<Vec<PositionStatusReport>> {
2863 let response = self
2865 .http_client
2866 .inner
2867 .get_subaccount(&self.wallet_address, self.subaccount_number)
2868 .await
2869 .context("failed to fetch subaccount from dYdX API")?;
2870
2871 let mut reports = Vec::new();
2872 let ts_init = UnixNanos::default();
2873
2874 for (market_ticker, perp_position) in &response.subaccount.open_perpetual_positions {
2875 let instrument = match self.get_instrument_by_market(market_ticker) {
2876 Some(inst) => inst,
2877 None => {
2878 log::warn!("Unknown market in position: {market_ticker}");
2879 continue;
2880 }
2881 };
2882
2883 if let Some(filter_id) = cmd.instrument_id
2884 && instrument.id() != filter_id
2885 {
2886 continue;
2887 }
2888
2889 let report = match parse_position_status_report(
2890 perp_position,
2891 &instrument,
2892 self.core.account_id,
2893 ts_init,
2894 ) {
2895 Ok(r) => r,
2896 Err(e) => {
2897 log::warn!("Failed to parse position status report: {e}");
2898 continue;
2899 }
2900 };
2901
2902 reports.push(report);
2903 }
2904
2905 Ok(reports)
2906 }
2907
2908 async fn generate_mass_status(
2909 &self,
2910 lookback_mins: Option<u64>,
2911 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
2912 let ts_init = UnixNanos::default();
2913
2914 let orders_response = self
2916 .http_client
2917 .inner
2918 .get_orders(
2919 &self.wallet_address,
2920 self.subaccount_number,
2921 None,
2922 Some(DYDX_INDEXER_REPORT_LIMIT),
2923 )
2924 .await
2925 .context("failed to fetch orders for mass status")?;
2926
2927 let subaccount_response = self
2929 .http_client
2930 .inner
2931 .get_subaccount(&self.wallet_address, self.subaccount_number)
2932 .await
2933 .context("failed to fetch subaccount for mass status")?;
2934
2935 let fills_response = self
2937 .http_client
2938 .inner
2939 .get_fills(
2940 &self.wallet_address,
2941 self.subaccount_number,
2942 None,
2943 Some(DYDX_INDEXER_REPORT_LIMIT),
2944 )
2945 .await
2946 .context("failed to fetch fills for mass status")?;
2947
2948 let mut order_reports = Vec::new();
2950 let mut orders_filtered = 0usize;
2951
2952 for order in orders_response {
2953 let instrument = match self.get_instrument_by_clob_pair_id(order.clob_pair_id) {
2954 Some(inst) => inst,
2955 None => {
2956 orders_filtered += 1;
2957 continue;
2958 }
2959 };
2960
2961 match parse_order_status_report(&order, &instrument, self.core.account_id, ts_init) {
2962 Ok(mut r) => {
2963 if !order.client_id.is_empty()
2964 && let Ok(client_id_u32) = order.client_id.parse::<u32>()
2965 {
2966 self.encoder.register_known_client_id(client_id_u32);
2967
2968 if let Some(decoded) = self
2969 .encoder
2970 .decode_if_known(client_id_u32, order.client_metadata)
2971 {
2972 log::debug!(
2973 "Decoded reconciliation order: dYdX client_id={} meta={:#x} -> '{}'",
2974 client_id_u32,
2975 order.client_metadata,
2976 decoded,
2977 );
2978 r.client_order_id = Some(decoded);
2979 }
2980 }
2981 order_reports.push(r);
2982 }
2983 Err(e) => {
2984 log::warn!("Failed to parse order status report: {e}");
2985 orders_filtered += 1;
2986 }
2987 }
2988 }
2989
2990 let mut position_reports = Vec::new();
2992
2993 for (market_ticker, perp_position) in
2994 &subaccount_response.subaccount.open_perpetual_positions
2995 {
2996 let instrument = match self.get_instrument_by_market(market_ticker) {
2997 Some(inst) => inst,
2998 None => continue,
2999 };
3000
3001 match parse_position_status_report(
3002 perp_position,
3003 &instrument,
3004 self.core.account_id,
3005 ts_init,
3006 ) {
3007 Ok(r) => position_reports.push(r),
3008 Err(e) => {
3009 log::warn!("Failed to parse position status report: {e}");
3010 }
3011 }
3012 }
3013
3014 let mut fill_reports = Vec::new();
3016 let mut fills_filtered = 0usize;
3017
3018 for fill in fills_response.fills {
3019 let instrument = match self.get_instrument_by_market(&fill.market) {
3020 Some(inst) => inst,
3021 None => {
3022 fills_filtered += 1;
3023 continue;
3024 }
3025 };
3026
3027 match parse_fill_report(&fill, &instrument, self.core.account_id, ts_init) {
3028 Ok(r) => fill_reports.push(r),
3029 Err(e) => {
3030 log::warn!("Failed to parse fill report: {e}");
3031 fills_filtered += 1;
3032 }
3033 }
3034 }
3035
3036 apply_avg_px_from_fills(&mut order_reports, &fill_reports);
3037
3038 {
3043 let cache = self.core.cache();
3044 order_reports.retain(|r| {
3045 let Some(cid) = r.client_order_id else {
3046 return true;
3047 };
3048
3049 match cache.order(&cid) {
3050 Some(order) if order.status().is_closed() => {
3051 log::debug!(
3052 "Skipping reconciliation report for terminal order {cid} \
3053 (cache status={:?}, venue status={:?})",
3054 order.status(),
3055 r.order_status,
3056 );
3057 false
3058 }
3059 _ => true,
3060 }
3061 });
3062 }
3063
3064 if let Some(mins) = lookback_mins {
3066 let now_ns = self.clock.get_time_ns();
3067 let cutoff_ns = now_ns.as_u64().saturating_sub(mins * 60 * 1_000_000_000);
3068 let cutoff = UnixNanos::from(cutoff_ns);
3069
3070 let orders_before = order_reports.len();
3071 order_reports.retain(|r| r.ts_last >= cutoff);
3072 let orders_removed = orders_before - order_reports.len();
3073
3074 let fills_before = fill_reports.len();
3075 fill_reports.retain(|r| r.ts_event >= cutoff);
3076 let fills_removed = fills_before - fill_reports.len();
3077
3078 log::debug!(
3079 "Lookback filter ({}min): orders {}->{} (removed {}), fills {}->{} (removed {}), positions {} (unfiltered)",
3080 mins,
3081 orders_before,
3082 order_reports.len(),
3083 orders_removed,
3084 fills_before,
3085 fill_reports.len(),
3086 fills_removed,
3087 position_reports.len(),
3088 );
3089 } else {
3090 log::debug!(
3091 "Generated mass status: {} orders ({} filtered), {} positions, {} fills ({} filtered)",
3092 order_reports.len(),
3093 orders_filtered,
3094 position_reports.len(),
3095 fill_reports.len(),
3096 fills_filtered,
3097 );
3098 }
3099
3100 let mut mass_status = ExecutionMassStatus::new(
3102 self.core.client_id,
3103 self.core.account_id,
3104 self.core.venue,
3105 ts_init,
3106 None, );
3108
3109 mass_status.add_order_reports(order_reports);
3110 mass_status.add_position_reports(position_reports);
3111 mass_status.add_fill_reports(fill_reports);
3112
3113 Ok(Some(mass_status))
3114 }
3115}
3116
3117struct PendingTaskLabel {
3118 id: u64,
3119 labels: Arc<Mutex<AHashMap<u64, &'static str>>>,
3120}
3121
3122impl Drop for PendingTaskLabel {
3123 fn drop(&mut self) {
3124 self.labels.lock().remove(&self.id);
3125 }
3126}
3127
3128#[allow(clippy::too_many_arguments)]
3132fn find_matching_order_report<F>(
3133 orders: &[crate::http::models::Order],
3134 instrument_filter: Option<InstrumentId>,
3135 client_order_id_filter: Option<ClientOrderId>,
3136 venue_order_id_filter: Option<VenueOrderId>,
3137 lookup_instrument: F,
3138 encoder: &ClientOrderIdEncoder,
3139 account_id: AccountId,
3140 ts_init: UnixNanos,
3141) -> anyhow::Result<Option<OrderStatusReport>>
3142where
3143 F: Fn(u32) -> Option<InstrumentAny>,
3144{
3145 for order in orders {
3146 let instrument = match lookup_instrument(order.clob_pair_id) {
3147 Some(inst) => inst,
3148 None => continue,
3149 };
3150
3151 if let Some(filter_id) = instrument_filter
3152 && instrument.id() != filter_id
3153 {
3154 continue;
3155 }
3156
3157 let mut report = parse_order_status_report(order, &instrument, account_id, ts_init)
3158 .context("failed to parse order status report")?;
3159
3160 if !order.client_id.is_empty()
3161 && let Ok(client_id_u32) = order.client_id.parse::<u32>()
3162 {
3163 encoder.register_known_client_id(client_id_u32);
3164
3165 if let Some(decoded) = encoder.decode_if_known(client_id_u32, order.client_metadata) {
3166 log::debug!(
3167 "Decoded order: dYdX client_id={} meta={:#x} -> '{}'",
3168 client_id_u32,
3169 order.client_metadata,
3170 decoded,
3171 );
3172 report.client_order_id = Some(decoded);
3173 }
3174 }
3175
3176 if let Some(client_order_id) = client_order_id_filter
3177 && report.client_order_id != Some(client_order_id)
3178 {
3179 continue;
3180 }
3181
3182 if let Some(venue_order_id) = venue_order_id_filter
3183 && report.venue_order_id.as_str() != venue_order_id.as_str()
3184 {
3185 continue;
3186 }
3187
3188 return Ok(Some(report));
3189 }
3190
3191 Ok(None)
3192}
3193
3194#[cfg(test)]
3195mod tests {
3196 use std::{cell::RefCell, rc::Rc};
3197
3198 use jiff::Timestamp;
3199 use nautilus_common::{
3200 cache::Cache, clock::TestClock, factories::OrderFactory, messages::ExecutionEvent,
3201 };
3202 use nautilus_model::{
3203 enums::OrderSide,
3204 identifiers::{Symbol, TraderId},
3205 instruments::{CryptoPerpetual, InstrumentAny},
3206 orders::{Order as _, OrderAny},
3207 types::{Currency, Price, Quantity},
3208 };
3209 use rstest::rstest;
3210 use rust_decimal_macros::dec;
3211
3212 use super::*;
3213 use crate::{
3214 common::{
3215 consts::DYDX_CLIENT_ID,
3216 enums::{DydxOrderStatus, DydxOrderType, DydxTimeInForce},
3217 },
3218 http::models::Order,
3219 };
3220
3221 fn test_instrument(symbol: &str, venue: &str) -> InstrumentAny {
3222 let instrument_id = InstrumentId::new(Symbol::new(symbol), Venue::new(venue));
3223 InstrumentAny::CryptoPerpetual(
3224 CryptoPerpetual::builder()
3225 .instrument_id(instrument_id)
3226 .raw_symbol(instrument_id.symbol)
3227 .base_currency(Currency::BTC())
3228 .quote_currency(Currency::USD())
3229 .settlement_currency(Currency::USD())
3230 .is_inverse(false)
3231 .price_precision(2)
3232 .size_precision(3)
3233 .price_increment(Price::new(0.01, 2))
3234 .size_increment(Quantity::new(0.001, 3))
3235 .ts_event(UnixNanos::default())
3236 .ts_init(UnixNanos::default())
3237 .build()
3238 .unwrap(),
3239 )
3240 }
3241
3242 fn test_order(id: &str, clob_pair_id: u32, client_id: &str) -> Order {
3243 Order {
3244 id: id.to_string(),
3245 subaccount_id: "sub-1".to_string(),
3246 client_id: client_id.to_string(),
3247 clob_pair_id,
3248 side: OrderSide::Buy,
3249 size: dec!(1.0),
3250 total_filled: dec!(0),
3251 price: dec!(50000),
3252 status: DydxOrderStatus::Open,
3253 order_type: DydxOrderType::Limit,
3254 time_in_force: DydxTimeInForce::Gtt,
3255 reduce_only: false,
3256 post_only: false,
3257 order_flags: 64,
3258 good_til_block: None,
3259 good_til_block_time: None,
3260 created_at_height: Some(100),
3261 client_metadata: 4,
3262 trigger_price: None,
3263 condition_type: None,
3264 conditional_order_trigger_subticks: None,
3265 execution: None,
3266 updated_at: None,
3267 updated_at_height: None,
3268 ticker: None,
3269 subaccount_number: 0,
3270 order_router_address: None,
3271 }
3272 }
3273
3274 fn test_order_factory() -> OrderFactory {
3275 let clock = Rc::new(RefCell::new(TestClock::new()));
3276 OrderFactory::new(
3277 TraderId::from("TRADER-001"),
3278 StrategyId::from("S-001"),
3279 Some(0),
3280 Some(0),
3281 clock,
3282 false,
3283 false,
3284 )
3285 }
3286
3287 fn create_execution_client() -> (
3288 DydxExecutionClient,
3289 Rc<RefCell<Cache>>,
3290 tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3291 ) {
3292 const TEST_PRIVATE_KEY: &str =
3293 "0000000000000000000000000000000000000000000000000000000000000001";
3294
3295 let cache = Rc::new(RefCell::new(Cache::default()));
3296 let core = ExecutionClientCore::new(
3297 TraderId::from("TRADER-001"),
3298 *DYDX_CLIENT_ID,
3299 *DYDX_VENUE,
3300 OmsType::Netting,
3301 AccountId::from("DYDX-001"),
3302 AccountType::Margin,
3303 None,
3304 cache.clone(),
3305 );
3306 core.set_connected();
3307
3308 let config = DydxAdapterConfig {
3309 base_url: "http://127.0.0.1:1".to_string(),
3310 ws_url: "ws://127.0.0.1:1".to_string(),
3311 grpc_url: "http://127.0.0.1:1".to_string(),
3312 grpc_urls: vec!["http://127.0.0.1:1".to_string()],
3313 wallet_address: Some("dydx1test".to_string()),
3314 private_key: Some(TEST_PRIVATE_KEY.to_string()),
3315 ..Default::default()
3316 };
3317
3318 let mut client =
3319 DydxExecutionClient::new(core, config, "dydx1test".to_string(), 0).unwrap();
3320 client
3321 .block_time_monitor
3322 .record_block(100, Timestamp::now());
3323
3324 let (sender, receiver) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
3325 client.emitter.set_sender(sender);
3326
3327 (client, cache, receiver)
3328 }
3329
3330 fn cache_order(cache: &Rc<RefCell<Cache>>, order: OrderAny) {
3331 cache
3332 .borrow_mut()
3333 .add_order(order, None, Some(*DYDX_CLIENT_ID), false)
3334 .unwrap();
3335 }
3336
3337 fn recv_order_event(
3338 rx: &mut tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3339 ) -> OrderEventAny {
3340 match rx.try_recv().expect("expected execution event") {
3341 ExecutionEvent::Order(event) => event,
3342 event => panic!("Expected order event, was {event:?}"),
3343 }
3344 }
3345
3346 #[rstest]
3347 fn test_submit_order_denies_quote_quantity() {
3348 let (client, cache, mut rx) = create_execution_client();
3349 let mut factory = test_order_factory();
3350 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3351 let order = factory.market(
3352 instrument_id,
3353 OrderSide::Buy,
3354 Quantity::from("100"),
3355 Some(TimeInForce::Ioc),
3356 None,
3357 Some(true),
3358 None,
3359 None,
3360 None,
3361 Some(ClientOrderId::from("O-QUOTE-QTY")),
3362 );
3363 cache_order(&cache, order.clone());
3364
3365 let command = SubmitOrder::from_order(
3366 &order,
3367 order.trader_id(),
3368 Some(*DYDX_CLIENT_ID),
3369 None,
3370 UUID4::new(),
3371 UnixNanos::default(),
3372 );
3373 client.submit_order(command).unwrap();
3374
3375 match recv_order_event(&mut rx) {
3376 OrderEventAny::Denied(event) => {
3377 assert_eq!(event.client_order_id, order.client_order_id());
3378 assert_eq!(event.instrument_id, instrument_id);
3379 assert_eq!(
3380 event.reason,
3381 "Quote quantity orders are not supported by dYdX"
3382 );
3383 }
3384 event => panic!("Expected order denied event, was {event:?}"),
3385 }
3386 }
3387
3388 #[rstest]
3389 fn test_submit_order_denies_market_fok() {
3390 let (client, cache, mut rx) = create_execution_client();
3391 let mut factory = test_order_factory();
3392 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3393 let order = factory.market(
3394 instrument_id,
3395 OrderSide::Buy,
3396 Quantity::from("0.1"),
3397 Some(TimeInForce::Fok),
3398 None,
3399 Some(false),
3400 None,
3401 None,
3402 None,
3403 Some(ClientOrderId::from("O-MARKET-FOK")),
3404 );
3405 cache_order(&cache, order.clone());
3406
3407 let command = SubmitOrder::from_order(
3408 &order,
3409 order.trader_id(),
3410 Some(*DYDX_CLIENT_ID),
3411 None,
3412 UUID4::new(),
3413 UnixNanos::default(),
3414 );
3415 client.submit_order(command).unwrap();
3416
3417 match recv_order_event(&mut rx) {
3418 OrderEventAny::Denied(event) => {
3419 assert_eq!(event.client_order_id, order.client_order_id());
3420 assert_eq!(event.instrument_id, instrument_id);
3421 assert!(
3422 event.reason.contains("Fill-or-kill") && event.reason.contains("deprecated"),
3423 "expected FOK-deprecation reason, was: {}",
3424 event.reason,
3425 );
3426 }
3427 event => panic!("Expected order denied event, was {event:?}"),
3428 }
3429 }
3430
3431 #[rstest]
3432 fn test_submit_order_denies_limit_fok() {
3433 let (client, cache, mut rx) = create_execution_client();
3437 let mut factory = test_order_factory();
3438 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3439 let order = factory.limit(
3440 instrument_id,
3441 OrderSide::Buy,
3442 Quantity::from("0.1"),
3443 Price::from("50000.0"),
3444 Some(TimeInForce::Fok),
3445 None,
3446 Some(false),
3447 None,
3448 Some(false),
3449 None,
3450 None,
3451 None,
3452 None,
3453 None,
3454 None,
3455 Some(ClientOrderId::from("O-LIMIT-FOK")),
3456 );
3457 cache_order(&cache, order.clone());
3458
3459 let command = SubmitOrder::from_order(
3460 &order,
3461 order.trader_id(),
3462 Some(*DYDX_CLIENT_ID),
3463 None,
3464 UUID4::new(),
3465 UnixNanos::default(),
3466 );
3467 client.submit_order(command).unwrap();
3468
3469 match recv_order_event(&mut rx) {
3470 OrderEventAny::Denied(event) => {
3471 assert_eq!(event.client_order_id, order.client_order_id());
3472 assert!(
3473 event.reason.contains("Fill-or-kill") && event.reason.contains("deprecated"),
3474 "expected FOK-deprecation reason, was: {}",
3475 event.reason,
3476 );
3477 }
3478 event => panic!("Expected order denied event, was {event:?}"),
3479 }
3480 }
3481
3482 #[rstest]
3483 fn test_submit_order_denies_day_tif() {
3484 let (client, cache, mut rx) = create_execution_client();
3487 let mut factory = test_order_factory();
3488 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3489 let order = factory.limit(
3490 instrument_id,
3491 OrderSide::Buy,
3492 Quantity::from("0.1"),
3493 Price::from("50000.0"),
3494 Some(TimeInForce::Day),
3495 None,
3496 Some(false),
3497 None,
3498 Some(false),
3499 None,
3500 None,
3501 None,
3502 None,
3503 None,
3504 None,
3505 Some(ClientOrderId::from("O-LIMIT-DAY")),
3506 );
3507 cache_order(&cache, order.clone());
3508
3509 let command = SubmitOrder::from_order(
3510 &order,
3511 order.trader_id(),
3512 Some(*DYDX_CLIENT_ID),
3513 None,
3514 UUID4::new(),
3515 UnixNanos::default(),
3516 );
3517 client.submit_order(command).unwrap();
3518
3519 match recv_order_event(&mut rx) {
3520 OrderEventAny::Denied(event) => {
3521 assert_eq!(event.client_order_id, order.client_order_id());
3522 assert!(
3523 event.reason.contains("DAY"),
3524 "expected DAY-unsupported reason, was: {}",
3525 event.reason,
3526 );
3527 }
3528 event => panic!("Expected order denied event, was {event:?}"),
3529 }
3530 }
3531
3532 #[rstest]
3536 #[case(TimeInForce::Fok, "Fill-or-kill")]
3537 #[case(TimeInForce::Day, "DAY")]
3538 fn test_submit_order_list_denies_unsupported_tif(
3539 #[case] tif: TimeInForce,
3540 #[case] reason_substring: &str,
3541 ) {
3542 use nautilus_common::messages::execution::SubmitOrderList;
3543 use nautilus_model::{identifiers::OrderListId, orders::OrderList};
3544
3545 let (client, cache, mut rx) = create_execution_client();
3546 let mut factory = test_order_factory();
3547 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3548 let order = factory.limit(
3549 instrument_id,
3550 OrderSide::Buy,
3551 Quantity::from("0.1"),
3552 Price::from("50000.0"),
3553 Some(tif),
3554 None,
3555 Some(false),
3556 None,
3557 Some(false),
3558 None,
3559 None,
3560 None,
3561 None,
3562 None,
3563 None,
3564 Some(ClientOrderId::from("O-LIST-1")),
3565 );
3566 cache_order(&cache, order.clone());
3567
3568 let order_list = OrderList::new(
3569 OrderListId::from("OL-1"),
3570 instrument_id,
3571 order.strategy_id(),
3572 vec![order.client_order_id()],
3573 UnixNanos::default(),
3574 );
3575 let init = order.init_event().clone();
3576
3577 let cmd = SubmitOrderList::new(
3578 order.trader_id(),
3579 Some(*DYDX_CLIENT_ID),
3580 order.strategy_id(),
3581 order_list,
3582 vec![init],
3583 None,
3584 None,
3585 None,
3586 UUID4::new(),
3587 UnixNanos::default(),
3588 None, );
3590 client.submit_order_list(cmd).unwrap();
3591
3592 match recv_order_event(&mut rx) {
3593 OrderEventAny::Denied(event) => {
3594 assert_eq!(event.client_order_id, order.client_order_id());
3595 assert!(
3596 event.reason.contains(reason_substring),
3597 "expected reason containing {reason_substring:?}, was: {}",
3598 event.reason,
3599 );
3600 }
3601 event => panic!("Expected order denied event, was {event:?}"),
3602 }
3603 }
3604
3605 #[rstest]
3608 fn test_submit_order_list_denies_supported_siblings_on_abort() {
3609 use nautilus_common::messages::execution::SubmitOrderList;
3610 use nautilus_model::{identifiers::OrderListId, orders::OrderList};
3611
3612 let (client, cache, mut rx) = create_execution_client();
3613 let mut factory = test_order_factory();
3614 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3615 let order_gtc = factory.limit(
3616 instrument_id,
3617 OrderSide::Buy,
3618 Quantity::from("0.1"),
3619 Price::from("50000.0"),
3620 Some(TimeInForce::Gtc),
3621 None,
3622 Some(false),
3623 None,
3624 Some(false),
3625 None,
3626 None,
3627 None,
3628 None,
3629 None,
3630 None,
3631 Some(ClientOrderId::from("O-MIXED-GTC")),
3632 );
3633 let order_fok = factory.limit(
3634 instrument_id,
3635 OrderSide::Sell,
3636 Quantity::from("0.1"),
3637 Price::from("51000.0"),
3638 Some(TimeInForce::Fok),
3639 None,
3640 Some(false),
3641 None,
3642 Some(false),
3643 None,
3644 None,
3645 None,
3646 None,
3647 None,
3648 None,
3649 Some(ClientOrderId::from("O-MIXED-FOK")),
3650 );
3651 cache_order(&cache, order_gtc.clone());
3652 cache_order(&cache, order_fok.clone());
3653
3654 let order_list = OrderList::new(
3655 OrderListId::from("OL-MIXED"),
3656 instrument_id,
3657 order_gtc.strategy_id(),
3658 vec![order_gtc.client_order_id(), order_fok.client_order_id()],
3659 UnixNanos::default(),
3660 );
3661
3662 let cmd = SubmitOrderList::new(
3663 order_gtc.trader_id(),
3664 Some(*DYDX_CLIENT_ID),
3665 order_gtc.strategy_id(),
3666 order_list,
3667 vec![
3668 order_gtc.init_event().clone(),
3669 order_fok.init_event().clone(),
3670 ],
3671 None,
3672 None,
3673 None,
3674 UUID4::new(),
3675 UnixNanos::default(),
3676 None, );
3678 client.submit_order_list(cmd).unwrap();
3679
3680 let mut denied_reasons = std::collections::HashMap::new();
3681
3682 for _ in 0..2 {
3683 match recv_order_event(&mut rx) {
3684 OrderEventAny::Denied(event) => {
3685 denied_reasons.insert(event.client_order_id, event.reason);
3686 }
3687 event => panic!("Expected order denied event, was {event:?}"),
3688 }
3689 }
3690
3691 assert!(
3692 denied_reasons[&order_gtc.client_order_id()].contains("sibling order"),
3693 "GTC sibling should be denied with list-abort reason",
3694 );
3695 assert!(
3696 denied_reasons[&order_fok.client_order_id()].contains("Fill-or-kill"),
3697 "FOK order should be denied with TIF reason",
3698 );
3699 }
3700
3701 #[tokio::test]
3704 async fn test_pending_task_labels_drain_with_scope() {
3705 let (client, _cache, _rx) = create_execution_client();
3706
3707 client.spawn_labeled("cancel_all_orders", async {
3708 tokio::time::sleep(std::time::Duration::from_secs(60)).await;
3709 });
3710 client.spawn_labeled("submit_order", async {
3711 tokio::time::sleep(std::time::Duration::from_secs(60)).await;
3712 });
3713
3714 {
3715 let labels = client.pending_task_labels.lock();
3716 assert_eq!(labels.len(), 2);
3717 }
3718
3719 client.begin_pending_shutdown();
3720 client
3721 .pending_tasks
3722 .finish_shutdown(Duration::ZERO, Duration::from_secs(1))
3723 .await
3724 .unwrap();
3725
3726 assert!(client.pending_tasks.is_empty());
3727 assert!(client.pending_task_labels.lock().is_empty());
3728 }
3729
3730 #[rstest]
3734 #[case("cancel_all_orders", true)]
3735 #[case("batch_cancel_orders", true)]
3736 #[case("Submit order", false)]
3737 #[case("batch_submit_short_term", false)]
3738 #[case("batch_submit_long_term", false)]
3739 fn test_pending_task_label_classification(
3740 #[case] label: &'static str,
3741 #[case] is_cancel: bool,
3742 ) {
3743 assert_eq!(label.contains("cancel"), is_cancel);
3747 }
3748
3749 #[rstest]
3750 fn test_emit_partitioned_cancel_rejections_emits_each_cancel() {
3751 let clock = get_atomic_clock_realtime();
3752 let mut emitter = ExecutionEventEmitter::new(
3753 clock,
3754 TraderId::from("TRADER-001"),
3755 AccountId::from("DYDX-001"),
3756 AccountType::Margin,
3757 None,
3758 );
3759 let (sender, mut rx) = tokio::sync::mpsc::unbounded_channel::<ExecutionEvent>();
3760 emitter.set_sender(sender);
3761
3762 let instrument_id = InstrumentId::from("BTC-USD-PERP.DYDX");
3763 let orders = vec![
3764 DydxCancelOrderRequest {
3765 instrument_id,
3766 client_id: 1,
3767 order_flags: types::ORDER_FLAG_SHORT_TERM,
3768 strategy_id: StrategyId::from("S-001"),
3769 client_order_id: ClientOrderId::from("O-CANCEL-1"),
3770 venue_order_id: Some(VenueOrderId::from("V-1")),
3771 },
3772 DydxCancelOrderRequest {
3773 instrument_id,
3774 client_id: 2,
3775 order_flags: types::ORDER_FLAG_LONG_TERM,
3776 strategy_id: StrategyId::from("S-001"),
3777 client_order_id: ClientOrderId::from("O-CANCEL-2"),
3778 venue_order_id: Some(VenueOrderId::from("V-2")),
3779 },
3780 ];
3781
3782 emit_partitioned_cancel_rejections(&orders, &emitter, clock, "broadcast failed");
3783
3784 for order in &orders {
3785 match recv_order_event(&mut rx) {
3786 OrderEventAny::CancelRejected(event) => {
3787 assert_eq!(event.client_order_id, order.client_order_id);
3788 assert_eq!(event.instrument_id, order.instrument_id);
3789 assert_eq!(event.venue_order_id, order.venue_order_id);
3790 assert_eq!(event.reason, "broadcast failed");
3791 }
3792 event => panic!("Expected order cancel rejected event, was {event:?}"),
3793 }
3794 }
3795 assert!(rx.try_recv().is_err());
3796 }
3797
3798 #[rstest]
3799 fn test_find_matching_order_report_returns_later_match() {
3800 let btc_inst = test_instrument("BTC-USD-PERP", "DYDX");
3804 let eth_inst = test_instrument("ETH-USD-PERP", "DYDX");
3805
3806 let orders = vec![
3808 test_order("order-eth", 1, "22222"),
3809 test_order("order-btc", 0, "11111"),
3810 ];
3811
3812 let encoder = ClientOrderIdEncoder::new();
3813 let report = find_matching_order_report(
3814 &orders,
3815 Some(btc_inst.id()),
3816 None,
3817 None,
3818 |clob_pair_id| match clob_pair_id {
3819 0 => Some(btc_inst.clone()),
3820 1 => Some(eth_inst.clone()),
3821 _ => None,
3822 },
3823 &encoder,
3824 AccountId::new("DYDX-001"),
3825 UnixNanos::default(),
3826 )
3827 .expect("lookup should succeed");
3828
3829 let report = report.expect("matching order should be found");
3830 assert_eq!(report.instrument_id, btc_inst.id());
3831 assert_eq!(report.venue_order_id.as_str(), "order-btc");
3832 }
3833
3834 #[rstest]
3835 fn test_find_matching_order_report_returns_none_when_no_match() {
3836 let btc_inst = test_instrument("BTC-USD-PERP", "DYDX");
3837 let eth_inst = test_instrument("ETH-USD-PERP", "DYDX");
3838
3839 let orders = vec![
3840 test_order("order-eth-1", 1, "22222"),
3841 test_order("order-eth-2", 1, "33333"),
3842 ];
3843
3844 let encoder = ClientOrderIdEncoder::new();
3845 let report = find_matching_order_report(
3846 &orders,
3847 Some(btc_inst.id()),
3848 None,
3849 None,
3850 |clob_pair_id| match clob_pair_id {
3851 0 => Some(btc_inst.clone()),
3852 1 => Some(eth_inst.clone()),
3853 _ => None,
3854 },
3855 &encoder,
3856 AccountId::new("DYDX-001"),
3857 UnixNanos::default(),
3858 )
3859 .expect("lookup should succeed");
3860
3861 assert!(report.is_none());
3862 }
3863
3864 #[rstest]
3865 fn test_find_matching_order_report_filters_by_venue_order_id() {
3866 let btc_inst = test_instrument("BTC-USD-PERP", "DYDX");
3867
3868 let orders = vec![
3869 test_order("order-a", 0, "11111"),
3870 test_order("order-b", 0, "22222"),
3871 test_order("order-c", 0, "33333"),
3872 ];
3873
3874 let encoder = ClientOrderIdEncoder::new();
3875 let target = VenueOrderId::new("order-b");
3876 let report = find_matching_order_report(
3877 &orders,
3878 None,
3879 None,
3880 Some(target),
3881 |_| Some(btc_inst.clone()),
3882 &encoder,
3883 AccountId::new("DYDX-001"),
3884 UnixNanos::default(),
3885 )
3886 .expect("lookup should succeed")
3887 .expect("matching order should be found");
3888
3889 assert_eq!(report.venue_order_id.as_str(), "order-b");
3890 }
3891
3892 #[rstest]
3893 fn test_find_matching_order_report_skips_orders_without_cached_instrument() {
3894 let btc_inst = test_instrument("BTC-USD-PERP", "DYDX");
3895
3896 let orders = vec![
3897 test_order("order-unknown", 99, "11111"),
3899 test_order("order-btc", 0, "22222"),
3900 ];
3901
3902 let encoder = ClientOrderIdEncoder::new();
3903 let report = find_matching_order_report(
3904 &orders,
3905 Some(btc_inst.id()),
3906 None,
3907 None,
3908 |clob_pair_id| (clob_pair_id == 0).then(|| btc_inst.clone()),
3909 &encoder,
3910 AccountId::new("DYDX-001"),
3911 UnixNanos::default(),
3912 )
3913 .expect("lookup should succeed")
3914 .expect("matching order should be found");
3915
3916 assert_eq!(report.venue_order_id.as_str(), "order-btc");
3917 }
3918}