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