Skip to main content

nautilus_dydx/execution/
mod.rs

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