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