Skip to main content

nautilus_derive/
execution.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution client implementation for the Derive adapter.
17//!
18//! Mirrors the Hyperliquid adapter's structural pattern: an
19//! [`ExecutionClientCore`] holds identity and connection state, an
20//! [`ExecutionEventEmitter`] publishes order/account events back to the live
21//! engine, and the venue clients ([`DeriveHttpClient`], [`DeriveWebSocketClient`])
22//! handle the wire. All state-changing requests are EIP-712 typed-data signed
23//! against the per-action module contracts on the Derive Chain; the
24//! `private/order` body in particular is built by [`order_to_derive_payload`].
25
26use std::{
27    sync::{
28        Arc,
29        atomic::{AtomicBool, Ordering},
30    },
31    time::{Duration, Instant},
32};
33
34use ahash::{AHashMap, AHashSet};
35use anyhow::Context;
36use async_trait::async_trait;
37use nautilus_common::{
38    cache::ORDER_NOT_FOUND,
39    clients::ExecutionClient,
40    live::runner::get_exec_event_sender,
41    messages::{
42        ExecutionReport,
43        execution::{
44            BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
45            GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
46            ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
47        },
48    },
49};
50use nautilus_core::{
51    AtomicMap, Params, UUID4, UnixNanos,
52    time::{AtomicTime, get_atomic_clock_realtime},
53};
54use nautilus_live::{
55    ExecutionClientCore, ExecutionEventEmitter, SocketControl,
56    task::{TaskGroup, TaskGroupGuard},
57};
58use nautilus_model::{
59    accounts::AccountAny,
60    data::QuoteTick,
61    enums::{OmsType, OrderSide, OrderStatus, OrderType, PositionSide},
62    events::{
63        OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderFilled, OrderRejected,
64    },
65    identifiers::{
66        AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, Symbol, Venue, VenueOrderId,
67    },
68    instruments::{Instrument, InstrumentAny},
69    orders::{Order, OrderAny},
70    reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
71    types::{AccountBalance, Currency, MarginBalance, Price, Quantity},
72};
73use rust_decimal::Decimal;
74use tokio_util::sync::CancellationToken;
75use ustr::Ustr;
76
77use crate::{
78    common::{
79        consts::{
80            DERIVE_ACCOUNT_REGISTRATION_TIMEOUT_SECS, DERIVE_VENUE, MIN_SIGNATURE_TTL,
81            TRIGGER_ORDER_SIGNATURE_TTL,
82        },
83        credential::DeriveCredential,
84        enums::{DeriveInstrumentType, DeriveOrderSide},
85        parse::{
86            derive_order_type_to_nautilus_for_order, derive_rejection_due_post_only,
87            format_instrument_id, format_venue_symbol,
88        },
89        retry::{http_retry_config, is_write_outcome_ambiguous_ws},
90    },
91    config::DeriveExecutionClientConfig,
92    http::{
93        DeriveCredentials, DeriveHttpClient,
94        models::{DeriveInstrument, DeriveOrder, DeriveReplaceOutcome, DeriveTrade},
95        parse::{
96            parse_derive_order_to_report_with_precision,
97            parse_derive_position_to_report_with_precision, parse_derive_subaccount_to_balances,
98            parse_derive_trade_to_fill_report_with_precision,
99        },
100        query::{
101            DeriveCancelByInstrumentParams, DeriveCancelByLabelParams, DeriveCancelParams,
102            DeriveCancelTriggerOrderParams, DeriveGetOpenOrdersParams, DeriveGetOrderHistoryParams,
103            DeriveGetOrderParams, DeriveGetPositionsParams, DeriveGetSubaccountParams,
104            DeriveGetTradeHistoryParams, DeriveGetTriggerOrdersParams,
105            order_replace_to_derive_payload, order_to_derive_payload,
106            trigger_order_to_derive_payload, validate_order_support,
107            validate_trigger_order_support,
108        },
109    },
110    signing::{
111        context::{SigningContext, resolve_signing_context},
112        nonce::{NonceError, NonceManager},
113    },
114    websocket::{
115        DeriveOrdersSubscriptionData, DeriveTradesSubscriptionData, DeriveWebSocketClient,
116        DeriveWsChannel, DeriveWsCredentials, DeriveWsError, DeriveWsExecutionHandle,
117        DeriveWsMessage, OrderIdentity, WsDispatchState, parse::parse_ticker_quote_from_rest,
118    },
119};
120
121const DERIVE_PRIVATE_PAGE_SIZE: u32 = 500;
122
123/// Live execution client for Derive.
124///
125/// Owns the HTTP and WebSocket clients used to talk to the venue plus an
126/// [`ExecutionEventEmitter`] that publishes order/account events back to the
127/// live engine. Order operations are signed against the per-environment
128/// EIP-712 signing context resolved at construction.
129#[derive(Debug)]
130pub struct DeriveExecutionClient {
131    core: ExecutionClientCore,
132    clock: &'static AtomicTime,
133    config: DeriveExecutionClientConfig,
134    credential: DeriveCredential,
135    emitter: ExecutionEventEmitter,
136    http_client: DeriveHttpClient,
137    ws_client: DeriveWebSocketClient,
138    ws_exec: DeriveWsExecutionHandle,
139    instruments: Arc<AtomicMap<InstrumentId, DeriveInstrument>>,
140    nonce_manager: Arc<NonceManager>,
141    signing: SigningContext,
142    is_connected: Arc<AtomicBool>,
143    cancellation_token: CancellationToken,
144    session_tasks: TaskGroup,
145    pending_tasks: TaskGroup,
146    shutdown_errors: Vec<String>,
147    dispatch_state: Arc<WsDispatchState>,
148}
149
150impl DeriveExecutionClient {
151    /// Creates a new [`DeriveExecutionClient`].
152    ///
153    /// Resolves wallet/session-key/subaccount from the supplied config, falling
154    /// back to the documented environment variables when fields are unset, and
155    /// parses the EIP-712 signing constants (domain separator, action typehash,
156    /// trade-module address) from config overrides or the shipped per-environment
157    /// defaults.
158    ///
159    /// # Errors
160    ///
161    /// Returns an error when:
162    /// - `max_fee_per_contract` is missing or not greater than zero.
163    /// - Required credentials are not provided via config or environment.
164    /// - Signing constants are still placeholders or cannot be parsed as hex.
165    /// - The HTTP or WebSocket client cannot be constructed.
166    pub fn new(
167        core: ExecutionClientCore,
168        config: DeriveExecutionClientConfig,
169    ) -> anyhow::Result<Self> {
170        config.validate()?;
171
172        let credential = DeriveCredential::resolve(
173            config.wallet_address.clone(),
174            config.session_key.clone(),
175            config.subaccount_id,
176            config.environment,
177        )?;
178
179        let http_credentials = DeriveCredentials::new(
180            credential.wallet_address().to_string(),
181            credential.session_key(),
182        )
183        .context("failed to build Derive HTTP credentials")?;
184        let retry_config = http_retry_config(
185            config.max_retries,
186            config.retry_delay_initial_ms,
187            config.retry_delay_max_ms,
188        );
189        let http_client = DeriveHttpClient::with_credentials(
190            config.rest_url(),
191            http_credentials,
192            Some(config.http_timeout_secs),
193            config.proxy_url.clone(),
194            Some(retry_config),
195        )
196        .context("failed to create Derive HTTP client")?;
197
198        let ws_credentials = DeriveWsCredentials::new(
199            credential.wallet_address().to_string(),
200            credential.session_key(),
201        )
202        .context("failed to build Derive WebSocket credentials")?;
203        let mut ws_client = DeriveWebSocketClient::with_credentials(
204            Some(config.ws_url()),
205            config.environment,
206            config.transport_backend,
207            config.proxy_url.clone(),
208            ws_credentials,
209            config.max_matching_requests_per_second,
210            config.max_per_instrument_matching_requests_per_second,
211        )
212        .with_socket_control(SocketControl::new(
213            core.client_id,
214            Some(*DERIVE_VENUE),
215            "derive-user-streams",
216        ));
217
218        if let Some(secs) = config.ws_timeout_secs {
219            ws_client.set_request_timeout(Duration::from_secs(secs));
220        }
221        // The handle shares the client's command channel, which survives the
222        // reconnect swap, so it stays valid for the client's lifetime.
223        let ws_exec = ws_client.execution_handle();
224
225        let signing = resolve_signing_context(&credential, &config)?;
226
227        let clock = get_atomic_clock_realtime();
228        let emitter = ExecutionEventEmitter::new(
229            clock,
230            core.trader_id,
231            core.account_id,
232            core.account_type,
233            core.base_currency,
234        );
235
236        let session_tasks = TaskGroup::new();
237        let pending_tasks = TaskGroup::new();
238
239        Ok(Self {
240            core,
241            clock,
242            config,
243            credential,
244            emitter,
245            http_client,
246            ws_client,
247            ws_exec,
248            instruments: Arc::new(AtomicMap::new()),
249            nonce_manager: Arc::new(NonceManager::new()),
250            signing,
251            is_connected: Arc::new(AtomicBool::new(false)),
252            cancellation_token: CancellationToken::new(),
253            session_tasks,
254            pending_tasks,
255            shutdown_errors: Vec::new(),
256            dispatch_state: Arc::new(WsDispatchState::new()),
257        })
258    }
259
260    /// Returns the resolved subaccount id.
261    #[must_use]
262    pub const fn subaccount_id(&self) -> u64 {
263        self.credential.subaccount_id()
264    }
265
266    /// Returns a reference to the resolved configuration.
267    #[must_use]
268    pub fn config(&self) -> &DeriveExecutionClientConfig {
269        &self.config
270    }
271
272    /// Returns a reference to the underlying HTTP client.
273    #[must_use]
274    pub fn http_client(&self) -> &DeriveHttpClient {
275        &self.http_client
276    }
277
278    /// Caches a Derive instrument by instrument ID so order submission can
279    /// resolve `base_asset_address` and `base_asset_sub_id` without
280    /// re-querying the venue.
281    pub fn cache_instrument(&self, instrument: DeriveInstrument) {
282        let instrument_id = format_instrument_id(instrument.instrument_name);
283        if let (Ok(price_increment), Ok(size_increment)) = (
284            Price::from_decimal(instrument.tick_size),
285            Quantity::from_decimal(instrument.amount_step),
286        ) {
287            self.dispatch_state.register_instrument_precision(
288                instrument_id,
289                price_increment.precision,
290                size_increment.precision,
291            );
292        }
293        self.instruments.insert(instrument_id, instrument);
294    }
295
296    /// Spawns a fire-and-forget task tracked in `pending_tasks` for teardown.
297    fn spawn_task<F>(&self, description: &'static str, fut: F)
298    where
299        F: std::future::Future<Output = anyhow::Result<()>> + Send + 'static,
300    {
301        let future = async move {
302            if let Err(e) = fut.await {
303                log::warn!("{description} failed: {e:?}");
304            }
305        };
306
307        if let Err(e) = self.pending_tasks.spawn(future) {
308            log::warn!("Skipping Derive {description} after shutdown began: {e}");
309        }
310    }
311
312    fn abort_pending_tasks(&self) {
313        self.pending_tasks.begin_shutdown();
314    }
315
316    fn abort_session_tasks(&self) {
317        self.session_tasks.begin_shutdown();
318        self.ws_client.begin_shutdown();
319    }
320
321    async fn await_pending_tasks(&self) -> anyhow::Result<()> {
322        self.pending_tasks.begin_shutdown();
323        self.pending_tasks
324            .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
325            .await
326            .map_err(|e| anyhow::anyhow!("Failed to terminate Derive execution tasks: {e}"))?;
327        Ok(())
328    }
329
330    async fn await_session_tasks(&self) -> anyhow::Result<()> {
331        self.session_tasks.begin_shutdown();
332        self.session_tasks
333            .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
334            .await
335            .map_err(|e| {
336                anyhow::anyhow!("Failed to terminate Derive execution session tasks: {e}")
337            })?;
338        Ok(())
339    }
340
341    async fn ensure_instruments_initialized(&self) -> anyhow::Result<()> {
342        if self.core.instruments_initialized() {
343            return Ok(());
344        }
345        // Lazy bootstrap: exec-side fetches per-instrument on first reference.
346        // Marking the flag prevents duplicate work across reconnect cycles.
347        self.core.set_instruments_initialized();
348        Ok(())
349    }
350
351    fn reconciliation_context(&self) -> DeriveReconciliationContext {
352        DeriveReconciliationContext {
353            http_client: self.http_client.clone(),
354            emitter: self.emitter.clone(),
355            client_id: self.core.client_id,
356            account_id: self.core.account_id,
357            subaccount_id: self.credential.subaccount_id(),
358            clock: self.clock,
359            dispatch_state: Arc::clone(&self.dispatch_state),
360        }
361    }
362
363    async fn refresh_account_state(&self) -> anyhow::Result<()> {
364        self.reconciliation_context().refresh_account_state().await
365    }
366
367    /// Blocks until the account appears in the cache, or `timeout_secs` elapses.
368    ///
369    /// The execution engine populates the cache from the [`refresh_account_state`]
370    /// event asynchronously; strategies that begin issuing orders before the
371    /// account is registered race the portfolio. Connecting blocks here so the
372    /// runner can rely on `core.cache().account(account_id)` immediately after
373    /// `connect()` returns.
374    async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
375        let account_id = self.core.account_id;
376
377        if self.core.cache().account(&account_id).is_some() {
378            log::info!("Account {account_id} registered");
379            return Ok(());
380        }
381
382        let start = Instant::now();
383        let timeout = Duration::from_secs_f64(timeout_secs);
384        let interval = Duration::from_millis(10);
385
386        loop {
387            tokio::time::sleep(interval).await;
388
389            if self.core.cache().account(&account_id).is_some() {
390                log::info!("Account {account_id} registered");
391                return Ok(());
392            }
393
394            if start.elapsed() >= timeout {
395                anyhow::bail!(
396                    "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
397                );
398            }
399        }
400    }
401
402    /// Reverses the partial state `connect()` set up before the failing step:
403    /// cancels the shared cancellation token, aborts the WS dispatch task,
404    /// and closes the WS client. Used when initial account state cannot be
405    /// loaded so that the next `connect()` call starts from a clean slate.
406    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
407        self.cancellation_token.cancel();
408        self.abort_session_tasks();
409        self.abort_pending_tasks();
410
411        if let Err(e) = self.ws_client.disconnect().await {
412            self.shutdown_errors
413                .push(format!("Derive WebSocket shutdown failed: {e}"));
414        }
415        let (session_result, pending_result) =
416            tokio::join!(self.await_session_tasks(), self.await_pending_tasks());
417        self.core.set_disconnected();
418        self.is_connected.store(false, Ordering::Release);
419
420        if let Err(e) = session_result {
421            self.shutdown_errors.push(e.to_string());
422        }
423
424        if let Err(e) = pending_result {
425            self.shutdown_errors.push(e.to_string());
426        }
427
428        if !self.shutdown_errors.is_empty() {
429            anyhow::bail!(std::mem::take(&mut self.shutdown_errors).join("; "));
430        }
431        Ok(())
432    }
433
434    fn start_ws_dispatch(
435        &self,
436        rx: tokio::sync::mpsc::UnboundedReceiver<DeriveWsMessage>,
437    ) -> anyhow::Result<()> {
438        let emitter = self.emitter.clone();
439        let account_id = self.core.account_id;
440        let clock = self.clock;
441        let cancellation = self.cancellation_token.clone();
442        let dispatch_state = self.dispatch_state.clone();
443        let reconciliation = self.reconciliation_context();
444        let is_connected = Arc::clone(&self.is_connected);
445        let session_spawner = self
446            .session_tasks
447            .spawner()
448            .map_err(|e| anyhow::anyhow!("Derive session task admission is closed: {e}"))?;
449
450        self.session_tasks.spawn(async move {
451            let mut rx = rx;
452
453            loop {
454                tokio::select! {
455                    biased;
456                    () = cancellation.cancelled() => break,
457                    maybe = rx.recv() => {
458                        match maybe {
459                            Some(DeriveWsMessage::Reconnected) => {
460                                let context = reconciliation.clone();
461                                let task_cancellation = cancellation.clone();
462
463                                if let Err(e) = session_spawner.spawn(async move {
464                                    tokio::select! {
465                                        () = task_cancellation.cancelled() => {}
466                                        result = context.recover_after_reconnect() => {
467                                            if let Err(e) = result {
468                                                log::warn!("Derive post-reconnect recovery failed: {e:?}");
469                                            }
470                                        }
471                                    }
472                                }) {
473                                    log::warn!("Skipping Derive reconnect recovery after shutdown began: {e}");
474                                }
475                            }
476                            Some(DeriveWsMessage::SessionRecoveryFailed(reason)) => {
477                                is_connected.store(false, Ordering::Release);
478                                log::error!("Derive execution WebSocket recovery failed: {reason}");
479                            }
480                            Some(DeriveWsMessage::Subscription(payload))
481                                if payload.channel.as_str().ends_with(".balances") =>
482                            {
483                                let context = reconciliation.clone();
484                                let task_cancellation = cancellation.clone();
485
486                                if let Err(e) = session_spawner.spawn(async move {
487                                    tokio::select! {
488                                        () = task_cancellation.cancelled() => {}
489                                        result = context.refresh_account_state() => {
490                                            if let Err(e) = result {
491                                                log::warn!("Derive balance update refresh failed: {e:?}");
492                                            }
493                                        }
494                                    }
495                                }) {
496                                    log::warn!("Skipping Derive account refresh after shutdown began: {e}");
497                                }
498                            }
499                            Some(message) => handle_ws_message(
500                                message,
501                                &emitter,
502                                account_id,
503                                clock,
504                                &dispatch_state,
505                            ),
506                            None => break,
507                        }
508                    }
509                }
510            }
511        })?;
512        Ok(())
513    }
514}
515
516#[async_trait(?Send)]
517impl ExecutionClient for DeriveExecutionClient {
518    fn is_connected(&self) -> bool {
519        self.is_connected.load(Ordering::Acquire)
520    }
521
522    fn client_id(&self) -> ClientId {
523        self.core.client_id
524    }
525
526    fn account_id(&self) -> AccountId {
527        self.core.account_id
528    }
529
530    fn venue(&self) -> Venue {
531        *DERIVE_VENUE
532    }
533
534    fn oms_type(&self) -> OmsType {
535        self.core.oms_type
536    }
537
538    fn get_account(&self) -> Option<AccountAny> {
539        self.core.cache().account_owned(&self.core.account_id)
540    }
541
542    fn start(&mut self) -> anyhow::Result<()> {
543        if self.core.is_started() {
544            return Ok(());
545        }
546
547        let sender = get_exec_event_sender();
548        self.emitter.set_sender(sender);
549        self.core.set_started();
550
551        log::info!(
552            "Started: client_id={}, account_id={}, subaccount_id={}, environment={:?}, proxy_url={:?}",
553            self.core.client_id,
554            self.core.account_id,
555            self.credential.subaccount_id(),
556            self.config.environment,
557            self.config.proxy_url,
558        );
559        Ok(())
560    }
561
562    fn stop(&mut self) -> anyhow::Result<()> {
563        if self.core.is_stopped() {
564            return Ok(());
565        }
566
567        log::info!("Stopping Derive execution client");
568
569        self.cancellation_token.cancel();
570        self.abort_session_tasks();
571        self.abort_pending_tasks();
572
573        self.core.set_stopped();
574        self.core.set_disconnected();
575        self.is_connected.store(false, Ordering::Release);
576
577        log::info!("Derive execution client stopped");
578        Ok(())
579    }
580
581    async fn connect(&mut self) -> anyhow::Result<()> {
582        if self.is_connected()
583            && !self.cancellation_token.is_cancelled()
584            && self.session_tasks.is_open()
585            && self.pending_tasks.is_open()
586        {
587            return Ok(());
588        }
589
590        log::info!("Connecting Derive execution client");
591
592        if self.cancellation_token.is_cancelled()
593            || !self.session_tasks.is_open()
594            || !self.pending_tasks.is_open()
595        {
596            self.teardown_partial_connect().await?;
597            self.session_tasks
598                .start_generation()
599                .map_err(|e| anyhow::anyhow!("Failed to start Derive session generation: {e}"))?;
600            self.pending_tasks
601                .start_generation()
602                .map_err(|e| anyhow::anyhow!("Failed to start Derive task generation: {e}"))?;
603            self.cancellation_token = CancellationToken::new();
604        }
605        let cancellation_token = self.cancellation_token.clone();
606        let ws_shutdown = self.ws_client.shutdown_handle();
607        let setup_guard =
608            TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
609                cancellation_token.cancel();
610                ws_shutdown.begin_shutdown();
611            });
612
613        self.ensure_instruments_initialized()
614            .await
615            .context("failed to initialize Derive instruments")?;
616
617        self.ws_client
618            .connect()
619            .await
620            .context("failed to connect Derive WebSocket")?;
621        let Some(rx) = self.ws_client.take_event_receiver() else {
622            let e = anyhow::anyhow!("Derive execution WS event receiver not initialized");
623            if let Err(teardown_error) = self.teardown_partial_connect().await {
624                return Err(e.context(format!(
625                    "Derive execution startup teardown failed: {teardown_error}"
626                )));
627            }
628            return Err(e);
629        };
630
631        let subaccount_id = self.credential.subaccount_id();
632        let channels = vec![
633            DeriveWsChannel::orders(subaccount_id),
634            DeriveWsChannel::private_trades(subaccount_id),
635            DeriveWsChannel::balances(subaccount_id),
636        ];
637
638        if let Err(e) = self.ws_client.subscribe_channels(channels).await {
639            log::warn!("Derive private WS subscriptions failed: {e}; tearing down");
640            if let Err(teardown_error) = self.teardown_partial_connect().await {
641                return Err(anyhow::Error::new(e).context(format!(
642                    "Derive execution startup teardown failed: {teardown_error}"
643                )));
644            }
645            return Err(anyhow::Error::new(e).context("failed Derive private WS subscriptions"));
646        }
647
648        if let Err(e) = self.start_ws_dispatch(rx) {
649            if let Err(teardown_error) = self.teardown_partial_connect().await {
650                return Err(e.context(format!(
651                    "Derive execution startup teardown failed: {teardown_error}"
652                )));
653            }
654            return Err(e.context("failed to register Derive execution WebSocket dispatch task"));
655        }
656
657        // Fail-fast if the initial account snapshot cannot load: without it,
658        // `await_account_registered` would block the full timeout window and
659        // surface a misleading registration timeout. Tear down the WS we
660        // already started so the caller does not leak the dispatch task.
661        if let Err(e) = self.refresh_account_state().await {
662            log::warn!("Initial Derive account state refresh failed: {e}; tearing down");
663            if let Err(teardown_error) = self.teardown_partial_connect().await {
664                return Err(e.context(format!(
665                    "Derive execution startup teardown failed: {teardown_error}"
666                )));
667            }
668            return Err(e.context("failed initial Derive account state refresh"));
669        }
670
671        if let Err(e) = self
672            .await_account_registered(DERIVE_ACCOUNT_REGISTRATION_TIMEOUT_SECS)
673            .await
674        {
675            log::warn!("Derive account did not register in time: {e}; tearing down");
676            if let Err(teardown_error) = self.teardown_partial_connect().await {
677                return Err(e.context(format!(
678                    "Derive execution startup teardown failed: {teardown_error}"
679                )));
680            }
681            return Err(e.context("failed waiting for Derive account registration"));
682        }
683
684        self.core.set_connected();
685        self.is_connected.store(true, Ordering::Release);
686        setup_guard.disarm();
687        log::info!(
688            "Connected Derive execution client ({:?})",
689            self.config.environment
690        );
691        Ok(())
692    }
693
694    async fn disconnect(&mut self) -> anyhow::Result<()> {
695        log::info!("Disconnecting Derive execution client");
696        self.teardown_partial_connect().await?;
697        log::info!("Derive execution client disconnected");
698        Ok(())
699    }
700
701    fn generate_account_state(
702        &self,
703        balances: Vec<AccountBalance>,
704        margins: Vec<MarginBalance>,
705        reported: bool,
706        ts_event: UnixNanos,
707        info: Option<Params>,
708    ) -> anyhow::Result<()> {
709        self.emitter
710            .emit_account_state(balances, margins, reported, ts_event, info);
711        Ok(())
712    }
713
714    fn on_instrument(&mut self, instrument: InstrumentAny) {
715        self.dispatch_state.register_instrument_precision(
716            instrument.id(),
717            instrument.price_precision(),
718            instrument.size_precision(),
719        );
720        // The exec-side instrument cache holds `DeriveInstrument` records so
721        // signing can pull `base_asset_address` / `base_asset_sub_id`; the
722        // generic `InstrumentAny` shape published on the bus does not carry
723        // those, so the data client populates the cache via
724        // [`Self::cache_instrument`] from its bootstrap pass instead.
725    }
726
727    async fn generate_order_status_report(
728        &self,
729        cmd: &GenerateOrderStatusReport,
730    ) -> anyhow::Result<Option<OrderStatusReport>> {
731        if cmd.venue_order_id.is_none() && cmd.client_order_id.is_none() {
732            log::warn!(
733                "Derive generate_order_status_report requires venue_order_id or client_order_id"
734            );
735            return Ok(None);
736        }
737
738        let subaccount_id = self.credential.subaccount_id();
739        let order = if let Some(venue_order_id) = cmd.venue_order_id {
740            match self
741                .http_client
742                .get_order(&DeriveGetOrderParams::new(
743                    subaccount_id,
744                    venue_order_id.as_str(),
745                ))
746                .await
747            {
748                Ok(order) => Some(order),
749                Err(e) => {
750                    let trigger_orders = self
751                        .http_client
752                        .get_trigger_orders(&DeriveGetTriggerOrdersParams::new(subaccount_id))
753                        .await?
754                        .orders;
755
756                    match trigger_orders
757                        .into_iter()
758                        .find(|o| o.order_id.as_str() == venue_order_id.as_str())
759                    {
760                        Some(order) => Some(order),
761                        None => return Err(e.into()),
762                    }
763                }
764            }
765        } else {
766            // Derive has no by-label lookup endpoint; scan open orders first,
767            // then trigger orders, then fall through to paginated history so
768            // terminal orders resolve for reconcilers that only carry the
769            // client_order_id.
770            let label = cmd.client_order_id.expect("guarded above");
771            let open_orders = self
772                .http_client
773                .get_open_orders(&DeriveGetOpenOrdersParams::new(subaccount_id))
774                .await?
775                .orders;
776            let mut found = open_orders
777                .into_iter()
778                .find(|o| o.label.as_str() == label.as_str());
779
780            if found.is_none() {
781                let trigger_orders = self
782                    .http_client
783                    .get_trigger_orders(&DeriveGetTriggerOrdersParams::new(subaccount_id))
784                    .await?
785                    .orders;
786                found = trigger_orders
787                    .into_iter()
788                    .find(|o| o.label.as_str() == label.as_str());
789            }
790
791            if found.is_none() {
792                let instrument_name = cmd.instrument_id.map(|id| id.symbol.as_str().to_string());
793                let mut page: u32 = 1;
794
795                'history: loop {
796                    let mut params = DeriveGetOrderHistoryParams::new(
797                        subaccount_id,
798                        page,
799                        DERIVE_PRIVATE_PAGE_SIZE,
800                    );
801
802                    if let Some(name) = instrument_name.as_deref() {
803                        params = params.with_instrument_name(name);
804                    }
805
806                    let result = self.http_client.get_order_history(&params).await?;
807                    let total_pages = result.pagination.num_pages;
808
809                    for order in result.orders {
810                        if order.label.as_str() == label.as_str() {
811                            found = Some(order);
812                            break 'history;
813                        }
814                    }
815
816                    if (page as i64) >= total_pages || total_pages == 0 {
817                        break;
818                    }
819                    page += 1;
820                }
821            }
822            found
823        };
824
825        let Some(order) = order else {
826            return Ok(None);
827        };
828
829        if let Some(instrument_id) = cmd.instrument_id
830            && InstrumentId::new(Symbol::new(order.instrument_name.as_str()), *DERIVE_VENUE)
831                != instrument_id
832        {
833            log::warn!(
834                "Derive order {} is for {} but report requested {}",
835                order.order_id,
836                order.instrument_name.as_str(),
837                instrument_id,
838            );
839            return Ok(None);
840        }
841
842        let (price_precision, size_precision) =
843            report_precision(&self.dispatch_state, order.instrument_name.as_str());
844        let ts_init = self.clock.get_time_ns();
845        let mut report = parse_derive_order_to_report_with_precision(
846            &order,
847            self.core.account_id,
848            price_precision,
849            size_precision,
850            ts_init,
851        )?;
852        // Prefer the parsed label (the venue's source of truth); only stamp
853        // the cmd's id when the venue order has no label at all.
854        if report.client_order_id.is_none()
855            && let Some(client_order_id) = cmd.client_order_id
856        {
857            report = report.with_client_order_id(client_order_id);
858        }
859        Ok(Some(report))
860    }
861
862    async fn generate_order_status_reports(
863        &self,
864        cmd: &GenerateOrderStatusReports,
865    ) -> anyhow::Result<Vec<OrderStatusReport>> {
866        self.reconciliation_context()
867            .generate_order_status_reports(cmd, false)
868            .await
869    }
870
871    async fn generate_fill_reports(
872        &self,
873        cmd: GenerateFillReports,
874    ) -> anyhow::Result<Vec<FillReport>> {
875        self.reconciliation_context()
876            .generate_fill_reports(cmd)
877            .await
878    }
879
880    async fn generate_position_status_reports(
881        &self,
882        cmd: &GeneratePositionStatusReports,
883    ) -> anyhow::Result<Vec<PositionStatusReport>> {
884        let snapshot = self
885            .reconciliation_context()
886            .generate_position_status_snapshot(cmd)
887            .await?;
888        Ok(snapshot.reports)
889    }
890
891    async fn generate_mass_status(
892        &self,
893        lookback_mins: Option<u64>,
894    ) -> anyhow::Result<Option<ExecutionMassStatus>> {
895        Box::pin(
896            self.reconciliation_context()
897                .generate_mass_status(lookback_mins),
898        )
899        .await
900        .map(Some)
901    }
902
903    fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
904        let order = self.core.cache().try_order_owned(&cmd.client_order_id)?;
905
906        if order.is_closed() {
907            log::warn!("Cannot submit closed order {}", order.client_order_id());
908            return Ok(());
909        }
910
911        // Deny before emit_order_submitted so unsupported fields never
912        // surface as venue rejections.
913        let is_trigger_order = is_derive_trigger_order_type(order.order_type());
914        let support = if is_trigger_order {
915            validate_trigger_order_support(&order)
916        } else {
917            validate_order_support(&order)
918        };
919
920        if let Err(e) = support {
921            let reason = e.to_string();
922            log::warn!("Cannot submit order {}: {reason}", order.client_order_id());
923            self.emitter.emit_order_denied(&order, &reason);
924            return Ok(());
925        }
926
927        // Spot has no position to reduce; the venue rejects reduce-only
928        // unconditionally (11025), so deny locally. Perp/option reduce-only is
929        // position-conditional and must still reach the venue.
930        if order.is_reduce_only()
931            && matches!(
932                self.core.cache().instrument(&cmd.instrument_id),
933                Some(InstrumentAny::CurrencyPair(_))
934            )
935        {
936            let reason = format!(
937                "reduce-only is not supported for spot instrument {}; Derive spot has no position to reduce",
938                cmd.instrument_id,
939            );
940            log::warn!("{reason}");
941            self.emitter.emit_order_denied(&order, &reason);
942            return Ok(());
943        }
944
945        // Keep the existing OrderDenied path here, then refresh before signing
946        let market_quote = if order.order_type() == OrderType::Market {
947            match self.core.cache().quote(&cmd.instrument_id) {
948                Some(_) => Some(()),
949                None => {
950                    let reason = format!(
951                        "no cached quote for {}; subscribe to quote data before submitting market orders",
952                        cmd.instrument_id,
953                    );
954                    log::warn!("{reason}");
955                    self.emitter.emit_order_denied(&order, &reason);
956                    return Ok(());
957                }
958            }
959        } else {
960            None
961        };
962
963        let venue_symbol = format_venue_symbol(&cmd.instrument_id)?.to_string();
964        let http_client = self.http_client.clone();
965        let ws_exec = self.ws_exec.clone();
966        let signing = self.signing.clone();
967        let nonce_manager = self.nonce_manager.clone();
968        let wallet_str = self.credential.wallet_address().to_string();
969        let emitter = self.emitter.clone();
970        let clock = self.clock;
971        let instruments = self.instruments.clone();
972        let instrument_id = cmd.instrument_id;
973        let order_for_task = order.clone();
974        let account_id = self.core.account_id;
975
976        // Capture identity so the WS dispatch can route subsequent updates
977        // for this order to proper events rather than execution reports.
978        let identity = OrderIdentity {
979            instrument_id: order.instrument_id(),
980            strategy_id: order.strategy_id(),
981            order_side: order.order_side(),
982            order_type: order.order_type(),
983        };
984        self.dispatch_state
985            .register_identity(order.client_order_id(), identity);
986
987        self.emitter.emit_order_submitted(&order);
988
989        let slippage_bps = self.signing.market_order_slippage_bps;
990        let dispatch_state = self.dispatch_state.clone();
991
992        self.spawn_task("submit_order", async move {
993            let instrument = match cached_or_fetch_instrument(
994                &http_client,
995                &instruments,
996                &instrument_id,
997                &venue_symbol,
998            )
999            .await
1000            {
1001                Ok(i) => i,
1002                Err(e) => {
1003                    log::warn!("Failed to resolve instrument {venue_symbol}: {e}");
1004                    dispatch_state.forget(&order_for_task.client_order_id());
1005                    let ts = clock.get_time_ns();
1006                    emitter.emit_order_rejected(
1007                        &order_for_task,
1008                        &format!("instrument resolution failed: {e}"),
1009                        ts,
1010                        false,
1011                    );
1012                    return Ok(());
1013                }
1014            };
1015
1016            // Lazy-resolution net: the synchronous deny is skipped when the
1017            // cache was empty at submit time. OrderSubmitted already fired, so
1018            // reject here rather than deny.
1019            if order_for_task.is_reduce_only()
1020                && instrument.instrument_type == DeriveInstrumentType::Erc20
1021            {
1022                let reason = format!(
1023                    "reduce-only is not supported for spot instrument {}; Derive spot has no position to reduce",
1024                    order_for_task.instrument_id(),
1025                );
1026                log::warn!("{reason}");
1027                dispatch_state.forget(&order_for_task.client_order_id());
1028                let ts = clock.get_time_ns();
1029                emitter.emit_order_rejected(&order_for_task, &reason, ts, false);
1030                return Ok(());
1031            }
1032
1033            // Avoid signing against a quote captured before instrument resolution
1034            let explicit_price = if market_quote.is_some() {
1035                let quote = match refresh_market_order_quote(
1036                    &http_client,
1037                    &venue_symbol,
1038                    &instrument,
1039                    clock,
1040                )
1041                .await
1042                {
1043                    Ok(quote) => quote,
1044                    Err(e) => {
1045                        let reason = format!(
1046                            "market-order quote refresh failed for {}: {e}",
1047                            order_for_task.client_order_id(),
1048                        );
1049                        log::warn!("{reason}");
1050                        dispatch_state.forget(&order_for_task.client_order_id());
1051                        let ts = clock.get_time_ns();
1052                        emitter.emit_order_rejected(&order_for_task, &reason, ts, false);
1053                        return Ok(());
1054                    }
1055                };
1056
1057                match market_order_limit_price(
1058                    &quote,
1059                    order_for_task.order_side(),
1060                    slippage_bps,
1061                    instrument.tick_size,
1062                ) {
1063                    Some(p) => Some(p),
1064                    None => {
1065                        let reason = format!(
1066                            "market-order slippage bound is non-positive for {} ({} bps)",
1067                            order_for_task.client_order_id(),
1068                            slippage_bps,
1069                        );
1070                        log::warn!("{reason}");
1071                        dispatch_state.forget(&order_for_task.client_order_id());
1072                        let ts = clock.get_time_ns();
1073                        emitter.emit_order_rejected(&order_for_task, &reason, ts, false);
1074                        return Ok(());
1075                    }
1076                }
1077            } else if matches!(
1078                order_for_task.order_type(),
1079                OrderType::StopMarket | OrderType::MarketIfTouched
1080            ) {
1081                let trigger_price = match order_for_task.trigger_price() {
1082                    Some(price) => price.as_decimal(),
1083                    None => {
1084                        let reason = format!(
1085                            "trigger market order {} is missing trigger_price",
1086                            order_for_task.client_order_id(),
1087                        );
1088                        log::warn!("{reason}");
1089                        dispatch_state.forget(&order_for_task.client_order_id());
1090                        let ts = clock.get_time_ns();
1091                        emitter.emit_order_rejected(&order_for_task, &reason, ts, false);
1092                        return Ok(());
1093                    }
1094                };
1095
1096                match trigger_market_limit_price(
1097                    trigger_price,
1098                    order_for_task.order_side(),
1099                    slippage_bps,
1100                    instrument.tick_size,
1101                ) {
1102                    Some(p) => Some(p),
1103                    None => {
1104                        let reason = format!(
1105                            "trigger market-order slippage bound is non-positive for {} ({} bps)",
1106                            order_for_task.client_order_id(),
1107                            slippage_bps,
1108                        );
1109                        log::warn!("{reason}");
1110                        dispatch_state.forget(&order_for_task.client_order_id());
1111                        let ts = clock.get_time_ns();
1112                        emitter.emit_order_rejected(&order_for_task, &reason, ts, false);
1113                        return Ok(());
1114                    }
1115                }
1116            } else {
1117                None
1118            };
1119
1120            let matching_reservation = match ws_exec
1121                .reserve_matching_request(
1122                    if is_trigger_order {
1123                        "private/trigger_order"
1124                    } else {
1125                        "private/order"
1126                    },
1127                    &instrument.instrument_name,
1128                )
1129                .await
1130            {
1131                Ok(reservation) => reservation,
1132                Err(e) => {
1133                    let (reason, due_post_only) = ws_rejection_reason(&e);
1134                    log::warn!(
1135                        "Cannot reserve Derive order quota for {}: {reason}",
1136                        order_for_task.client_order_id(),
1137                    );
1138                    dispatch_state.forget(&order_for_task.client_order_id());
1139                    let ts = clock.get_time_ns();
1140                    emitter.emit_order_rejected(
1141                        &order_for_task,
1142                        &reason,
1143                        ts,
1144                        due_post_only,
1145                    );
1146                    return Ok(());
1147                }
1148            };
1149
1150            if is_trigger_order {
1151                let nonce = match resolve_submit_nonce(
1152                    nonce_manager.next_nonce(&wallet_str, signing.subaccount_id),
1153                    &emitter,
1154                    &dispatch_state,
1155                    &order_for_task,
1156                    clock,
1157                ) {
1158                    Some(nonce) => nonce,
1159                    None => return Ok(()),
1160                };
1161                let expiry = trigger_order_signature_expiry(clock);
1162                let payload = match trigger_order_to_derive_payload(
1163                    &order_for_task,
1164                    &instrument,
1165                    signing.subaccount_id,
1166                    signing.wallet_address,
1167                    &signing.signer,
1168                    nonce,
1169                    expiry,
1170                    signing.trade_module_address,
1171                    signing.domain_separator,
1172                    signing.action_typehash,
1173                    signing.max_fee_per_contract,
1174                    explicit_price,
1175                    ws_exec.conn_id(),
1176                    UUID4::new().to_string(),
1177                ) {
1178                    Ok(p) => p,
1179                    Err(e) => {
1180                        log::warn!(
1181                            "Trigger order encode failed for {}: {e}",
1182                            order_for_task.client_order_id()
1183                        );
1184                        dispatch_state.forget(&order_for_task.client_order_id());
1185                        let ts = clock.get_time_ns();
1186                        emitter.emit_order_rejected(
1187                            &order_for_task,
1188                            &format!("order encoding failed: {e}"),
1189                            ts,
1190                            false,
1191                        );
1192                        return Ok(());
1193                    }
1194                };
1195
1196                log::debug!(
1197                    "Derive trigger submit payload client_order_id={} instrument_name={} direction={} order_type={} time_in_force={} amount={} limit_price={} trigger_price={:?} trigger_price_type={:?} trigger_type={:?}",
1198                    order_for_task.client_order_id(),
1199                    payload.order.instrument_name.as_str(),
1200                    payload.order.direction,
1201                    payload.order.order_type,
1202                    payload.order.time_in_force,
1203                    payload.order.amount,
1204                    payload.order.limit_price,
1205                    payload.order.trigger_price,
1206                    payload.order.trigger_price_type,
1207                    payload.order.trigger_type,
1208                );
1209
1210                match ws_exec
1211                    .submit_trigger_order_after_rate_limit(&payload, matching_reservation)
1212                    .await
1213                {
1214                    Ok(order) => {
1215                        let venue_order_id = VenueOrderId::new(order.order_id.as_str());
1216                        dispatch_state.record_venue_order_id(
1217                            order_for_task.client_order_id(),
1218                            venue_order_id,
1219                        );
1220                        let ts_now = clock.get_time_ns();
1221                        ensure_accepted_emitted(
1222                            &emitter,
1223                            &dispatch_state,
1224                            order_for_task.client_order_id(),
1225                            identity,
1226                            venue_order_id,
1227                            account_id,
1228                            ts_now,
1229                            ts_now,
1230                        );
1231                        log::debug!(
1232                            "Trigger order submitted: client_order_id={} venue_order_id={venue_order_id}",
1233                            order_for_task.client_order_id(),
1234                        );
1235                    }
1236                    Err(e) if is_write_outcome_ambiguous_ws(&e) => {
1237                        log::warn!(
1238                            "Derive trigger submit for {} returned ambiguous WS outcome: {e}; awaiting reconciliation",
1239                            order_for_task.client_order_id(),
1240                        );
1241                    }
1242                    Err(e) => {
1243                        let (reason, due_post_only) = ws_rejection_reason(&e);
1244                        log::debug!(
1245                            "Derive rejected trigger order {}: {reason}",
1246                            order_for_task.client_order_id(),
1247                        );
1248                        dispatch_state.forget(&order_for_task.client_order_id());
1249                        let ts = clock.get_time_ns();
1250                        emitter.emit_order_rejected(
1251                            &order_for_task,
1252                            &reason,
1253                            ts,
1254                            due_post_only,
1255                        );
1256                    }
1257                }
1258                return Ok(());
1259            }
1260
1261            let expiry =
1262                match normal_order_signature_expiry(clock, signing.signature_expiry_secs) {
1263                    Ok(expiry) => expiry,
1264                    Err(e) => {
1265                        log::warn!(
1266                            "Order expiry validation failed for {}: {e}",
1267                            order_for_task.client_order_id()
1268                        );
1269                        dispatch_state.forget(&order_for_task.client_order_id());
1270                        let ts = clock.get_time_ns();
1271                        emitter.emit_order_rejected(
1272                            &order_for_task,
1273                            &format!("order expiry validation failed: {e}"),
1274                            ts,
1275                            false,
1276                        );
1277                        return Ok(());
1278                    }
1279                };
1280            let nonce = match resolve_submit_nonce(
1281                nonce_manager.next_nonce(&wallet_str, signing.subaccount_id),
1282                &emitter,
1283                &dispatch_state,
1284                &order_for_task,
1285                clock,
1286            ) {
1287                Some(nonce) => nonce,
1288                None => return Ok(()),
1289            };
1290            let payload = match order_to_derive_payload(
1291                &order_for_task,
1292                &instrument,
1293                signing.subaccount_id,
1294                signing.wallet_address,
1295                &signing.signer,
1296                nonce,
1297                expiry,
1298                signing.trade_module_address,
1299                signing.domain_separator,
1300                signing.action_typehash,
1301                signing.max_fee_per_contract,
1302                explicit_price,
1303            ) {
1304                Ok(p) => p,
1305                Err(e) => {
1306                    log::warn!("Order encode failed for {}: {e}", order_for_task.client_order_id());
1307                    dispatch_state.forget(&order_for_task.client_order_id());
1308                    let ts = clock.get_time_ns();
1309                    emitter.emit_order_rejected(
1310                        &order_for_task,
1311                        &format!("order encoding failed: {e}"),
1312                        ts,
1313                        false,
1314                    );
1315                    return Ok(());
1316                }
1317            };
1318
1319            // Pre-flight debug log so a venue 11012-style rejection can be
1320            // diagnosed without re-running with full payload tracing.
1321            log::debug!(
1322                "Derive submit payload client_order_id={} instrument_name={} direction={} order_type={} time_in_force={} amount={} limit_price={}",
1323                order_for_task.client_order_id(),
1324                payload.instrument_name.as_str(),
1325                payload.direction,
1326                payload.order_type,
1327                payload.time_in_force,
1328                payload.amount,
1329                payload.limit_price,
1330            );
1331
1332            // Discard the result (and any `trades` it carries): fills arrive on
1333            // the `.trades` channel and are deduped by trade id.
1334            match ws_exec
1335                .submit_order_after_rate_limit(&payload, matching_reservation)
1336                .await
1337            {
1338                Ok(_) => {
1339                    log::debug!(
1340                        "Order submitted: client_order_id={}",
1341                        order_for_task.client_order_id(),
1342                    );
1343                }
1344                // See docs/integrations/derive.md "Order rejection semantics".
1345                Err(e) if is_write_outcome_ambiguous_ws(&e) => {
1346                    log::warn!(
1347                        "Derive submit for {} returned ambiguous WS outcome: {e}; awaiting reconciliation",
1348                        order_for_task.client_order_id(),
1349                    );
1350                }
1351                Err(e) => {
1352                    let (reason, due_post_only) = ws_rejection_reason(&e);
1353                    log::debug!(
1354                        "Derive rejected order {}: {reason}",
1355                        order_for_task.client_order_id(),
1356                    );
1357                    dispatch_state.forget(&order_for_task.client_order_id());
1358                    let ts = clock.get_time_ns();
1359                    emitter.emit_order_rejected(&order_for_task, &reason, ts, due_post_only);
1360                }
1361            }
1362            Ok(())
1363        });
1364
1365        Ok(())
1366    }
1367
1368    fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
1369        let orders = self.core.get_orders_for_list(&cmd.order_list)?;
1370        for order in orders {
1371            let sub = SubmitOrder::from_order(
1372                &order,
1373                cmd.trader_id,
1374                cmd.client_id,
1375                cmd.position_id,
1376                UUID4::new(),
1377                cmd.ts_init,
1378            );
1379            self.submit_order(sub)?;
1380        }
1381        Ok(())
1382    }
1383
1384    fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
1385        let http_client = self.http_client.clone();
1386        let ws_exec = self.ws_exec.clone();
1387        let subaccount_id = self.credential.subaccount_id();
1388        let venue_symbol = format_venue_symbol(&cmd.instrument_id)?.to_string();
1389        let emitter = self.emitter.clone();
1390        let clock = self.clock;
1391        let account_id = self.core.account_id;
1392        let dispatch_state = self.dispatch_state.clone();
1393        let strategy_id = cmd.strategy_id;
1394        let instrument_id = cmd.instrument_id;
1395        let client_order_id = cmd.client_order_id;
1396        let venue_order_id = cmd.venue_order_id;
1397        let is_trigger_order = self
1398            .core
1399            .cache()
1400            .order(&client_order_id)
1401            .is_some_and(|order| is_derive_trigger_order_type(order.order_type()));
1402
1403        self.spawn_task("cancel_order", async move {
1404            let outcome = match venue_order_id {
1405                Some(venue_order_id) if is_trigger_order => {
1406                    ws_exec
1407                        .cancel_trigger_order(&DeriveCancelTriggerOrderParams::new(
1408                            subaccount_id,
1409                            venue_order_id.as_str(),
1410                        ))
1411                        .await
1412                        .map(Some)
1413                }
1414                Some(venue_order_id) => ws_exec
1415                    .cancel_order(&DeriveCancelParams::new(
1416                        subaccount_id,
1417                        venue_symbol.as_str(),
1418                        venue_order_id.as_str(),
1419                    ))
1420                    .await
1421                    .map(|()| None),
1422                None if is_trigger_order => {
1423                    let trigger_orders = match http_client
1424                        .get_trigger_orders(&DeriveGetTriggerOrdersParams::new(subaccount_id))
1425                        .await
1426                    {
1427                        Ok(result) => result.orders,
1428                        Err(e) => {
1429                            let reason = format!("failed to resolve trigger order by label: {e}");
1430                            log::warn!("Cannot cancel trigger order {client_order_id}: {reason}");
1431                            emitter.emit_order_cancel_rejected_event(
1432                                strategy_id,
1433                                instrument_id,
1434                                client_order_id,
1435                                None,
1436                                &reason,
1437                                clock.get_time_ns(),
1438                            );
1439                            return Ok(());
1440                        }
1441                    };
1442                    let Some(trigger_order) = trigger_orders.into_iter().find(|order| {
1443                        order.label.as_str() == client_order_id.as_str()
1444                            && order.instrument_name.as_str() == venue_symbol
1445                    }) else {
1446                        let reason = "trigger order not found for client_order_id";
1447                        log::warn!("Cannot cancel trigger order {client_order_id}: {reason}");
1448                        emitter.emit_order_cancel_rejected_event(
1449                            strategy_id,
1450                            instrument_id,
1451                            client_order_id,
1452                            None,
1453                            reason,
1454                            clock.get_time_ns(),
1455                        );
1456                        return Ok(());
1457                    };
1458                    ws_exec
1459                        .cancel_trigger_order(&DeriveCancelTriggerOrderParams::new(
1460                            subaccount_id,
1461                            trigger_order.order_id.as_str(),
1462                        ))
1463                        .await
1464                        .map(Some)
1465                }
1466                None => ws_exec
1467                    .cancel_by_label(&DeriveCancelByLabelParams::new(
1468                        subaccount_id,
1469                        client_order_id.as_str(),
1470                    ))
1471                    .await
1472                    .map(|result| {
1473                        if result.cancelled_orders == 0 {
1474                            let reason = "no open order matched the client_order_id label";
1475                            log::debug!(
1476                                "Derive rejected cancel for {client_order_id}: {reason}"
1477                            );
1478                            let ts = clock.get_time_ns();
1479                            emitter.emit_order_cancel_rejected_event(
1480                                strategy_id,
1481                                instrument_id,
1482                                client_order_id,
1483                                None,
1484                                reason,
1485                                ts,
1486                            );
1487                        }
1488                        None
1489                    }),
1490            };
1491
1492            match outcome {
1493                Ok(Some(canceled_order)) => {
1494                    let canceled_venue_order_id =
1495                        VenueOrderId::new(canceled_order.order_id.as_str());
1496                    let ts = clock.get_time_ns();
1497
1498                    ensure_canceled_emitted(
1499                        &emitter,
1500                        &dispatch_state,
1501                        client_order_id,
1502                        OrderIdentity {
1503                            instrument_id,
1504                            strategy_id,
1505                            order_side: match canceled_order.direction {
1506                                DeriveOrderSide::Buy => OrderSide::Buy,
1507                                DeriveOrderSide::Sell => OrderSide::Sell,
1508                            },
1509                            order_type: derive_order_type_to_nautilus_for_order(
1510                                canceled_order.order_type,
1511                                canceled_order.trigger_type,
1512                            ),
1513                        },
1514                        canceled_venue_order_id,
1515                        account_id,
1516                        ts,
1517                        ts,
1518                    );
1519                    dispatch_state.forget(&client_order_id);
1520                }
1521                Ok(None) => {}
1522                // See docs/integrations/derive.md "Order rejection semantics".
1523                Err(e) if is_write_outcome_ambiguous_ws(&e) => {
1524                    log::warn!(
1525                        "Derive cancel for {client_order_id} returned ambiguous WS outcome: {e}; awaiting reconciliation",
1526                    );
1527                }
1528                Err(e) => {
1529                    let (reason, _) = ws_rejection_reason(&e);
1530                    log::debug!("Derive rejected cancel for {client_order_id}: {reason}");
1531                    let ts = clock.get_time_ns();
1532                    emitter.emit_order_cancel_rejected_event(
1533                        strategy_id,
1534                        instrument_id,
1535                        client_order_id,
1536                        venue_order_id,
1537                        &reason,
1538                        ts,
1539                    );
1540                }
1541            }
1542            Ok(())
1543        });
1544        Ok(())
1545    }
1546
1547    fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
1548        let venue_symbol = format_venue_symbol(&cmd.instrument_id)?.to_string();
1549        let side_filter = cmd.order_side;
1550        let cache = self.core.cache();
1551        let orders = cache.orders_open_refs(
1552            Some(&self.core.venue),
1553            Some(&cmd.instrument_id),
1554            None,
1555            Some(&self.core.account_id),
1556            side_filter,
1557        );
1558        let mut cancels = Vec::with_capacity(orders.len());
1559
1560        for order in orders {
1561            let client_order_id = order.client_order_id();
1562            if cache.client_id(&client_order_id) != Some(&self.core.client_id) {
1563                continue;
1564            }
1565
1566            let is_trigger = is_derive_trigger_order_type(order.order_type());
1567            if side_filter.is_none() && !is_trigger {
1568                continue;
1569            }
1570
1571            let Some(venue_order_id) = order.venue_order_id() else {
1572                log::warn!(
1573                    "Cannot cancel all orders for {}: order {client_order_id} has no venue_order_id",
1574                    cmd.instrument_id,
1575                );
1576                return Ok(());
1577            };
1578            cancels.push((venue_order_id, is_trigger));
1579        }
1580        drop(cache);
1581
1582        if side_filter.is_some() && cancels.is_empty() {
1583            return Ok(());
1584        }
1585
1586        let ws_exec = self.ws_exec.clone();
1587        let subaccount_id = self.credential.subaccount_id();
1588
1589        self.spawn_task("cancel_all_orders", async move {
1590            for (venue_order_id, is_trigger) in cancels {
1591                let outcome = if is_trigger {
1592                    ws_exec
1593                        .cancel_trigger_order(&DeriveCancelTriggerOrderParams::new(
1594                            subaccount_id,
1595                            venue_order_id.as_str(),
1596                        ))
1597                        .await
1598                        .map(|_| ())
1599                } else {
1600                    ws_exec
1601                        .cancel_order(&DeriveCancelParams::new(
1602                            subaccount_id,
1603                            venue_symbol.as_str(),
1604                            venue_order_id.as_str(),
1605                        ))
1606                        .await
1607                };
1608
1609                if let Err(e) = outcome {
1610                    log::warn!(
1611                        "Derive cancel_all_orders: cancel for {venue_order_id} failed: {e}",
1612                    );
1613                }
1614            }
1615
1616            if side_filter.is_none() {
1617                match ws_exec
1618                    .cancel_by_instrument(&DeriveCancelByInstrumentParams::new(
1619                        subaccount_id,
1620                        venue_symbol.as_str(),
1621                    ))
1622                    .await
1623                {
1624                    Ok(result) if result.cancelled_orders == 0 => {
1625                        log::debug!("No open orders to cancel for {venue_symbol}");
1626                    }
1627                    Ok(_) => {}
1628                    Err(e) => {
1629                        log::warn!("Derive cancel_all_orders failed for {venue_symbol}: {e}");
1630                    }
1631                }
1632            }
1633            Ok(())
1634        });
1635        Ok(())
1636    }
1637
1638    fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1639        for inner in cmd.cancels {
1640            self.cancel_order(inner)?;
1641        }
1642        Ok(())
1643    }
1644
1645    fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
1646        let ts_now = self.clock.get_time_ns();
1647
1648        let Some(venue_order_id) = cmd.venue_order_id else {
1649            let reason = "venue_order_id is required for modify";
1650            log::warn!("Cannot modify order {}: {reason}", cmd.client_order_id);
1651            self.emitter.emit_order_modify_rejected_event(
1652                cmd.strategy_id,
1653                cmd.instrument_id,
1654                cmd.client_order_id,
1655                None,
1656                reason,
1657                ts_now,
1658            );
1659            return Ok(());
1660        };
1661
1662        let Ok(order) = self.core.cache().try_order_owned(&cmd.client_order_id) else {
1663            let reason = ORDER_NOT_FOUND;
1664            log::warn!("Cannot modify order {}: {reason}", cmd.client_order_id);
1665            self.emitter.emit_order_modify_rejected_event(
1666                cmd.strategy_id,
1667                cmd.instrument_id,
1668                cmd.client_order_id,
1669                Some(venue_order_id),
1670                reason,
1671                ts_now,
1672            );
1673            return Ok(());
1674        };
1675
1676        if is_derive_trigger_order_type(order.order_type()) {
1677            let reason = "Derive trigger orders cannot be modified; cancel and resubmit";
1678            log::warn!("Cannot modify order {}: {reason}", cmd.client_order_id);
1679            self.emitter.emit_order_modify_rejected_event(
1680                cmd.strategy_id,
1681                cmd.instrument_id,
1682                cmd.client_order_id,
1683                Some(venue_order_id),
1684                reason,
1685                ts_now,
1686            );
1687            return Ok(());
1688        }
1689
1690        let target_quantity = cmd.quantity.unwrap_or_else(|| order.quantity());
1691        let target_price = cmd.price.or_else(|| order.price());
1692
1693        let venue_symbol = format_venue_symbol(&cmd.instrument_id)?.to_string();
1694        let http_client = self.http_client.clone();
1695        let ws_exec = self.ws_exec.clone();
1696        let signing = self.signing.clone();
1697        let nonce_manager = self.nonce_manager.clone();
1698        let wallet_str = self.credential.wallet_address().to_string();
1699        let emitter = self.emitter.clone();
1700        let clock = self.clock;
1701        let instruments = self.instruments.clone();
1702        let dispatch_state = self.dispatch_state.clone();
1703        let order_for_task = order;
1704        let strategy_id = cmd.strategy_id;
1705        let instrument_id = cmd.instrument_id;
1706        let client_order_id = cmd.client_order_id;
1707        let stale_venue_order_id = venue_order_id;
1708        let account_id = self.core.account_id;
1709        let voi_str = venue_order_id.to_string();
1710
1711        self.spawn_task("modify_order", async move {
1712            let instrument = match cached_or_fetch_instrument(
1713                &http_client,
1714                &instruments,
1715                &instrument_id,
1716                &venue_symbol,
1717            )
1718            .await
1719            {
1720                Ok(i) => i,
1721                Err(e) => {
1722                    let reason = format!("instrument resolution failed: {e}");
1723                    log::warn!("Cannot modify order {client_order_id}: {reason}");
1724                    let ts = clock.get_time_ns();
1725                    emitter.emit_order_modify_rejected_event(
1726                        strategy_id,
1727                        instrument_id,
1728                        client_order_id,
1729                        Some(stale_venue_order_id),
1730                        &reason,
1731                        ts,
1732                    );
1733                    return Ok(());
1734                }
1735            };
1736
1737            let matching_reservation = match ws_exec
1738                .reserve_matching_request("private/replace", &instrument.instrument_name)
1739                .await
1740            {
1741                Ok(reservation) => reservation,
1742                Err(e) => {
1743                    let (reason, _) = ws_rejection_reason(&e);
1744                    log::warn!("Cannot reserve Derive replace quota for {client_order_id}: {reason}");
1745                    let ts = clock.get_time_ns();
1746                    emitter.emit_order_modify_rejected_event(
1747                        strategy_id,
1748                        instrument_id,
1749                        client_order_id,
1750                        Some(stale_venue_order_id),
1751                        &reason,
1752                        ts,
1753                    );
1754                    return Ok(());
1755                }
1756            };
1757
1758            let expiry = match normal_order_signature_expiry(clock, signing.signature_expiry_secs) {
1759                Ok(expiry) => expiry,
1760                Err(e) => {
1761                    let reason = format!("replace expiry validation failed: {e}");
1762                    log::warn!("Cannot modify order {client_order_id}: {reason}");
1763                    let ts = clock.get_time_ns();
1764                    emitter.emit_order_modify_rejected_event(
1765                        strategy_id,
1766                        instrument_id,
1767                        client_order_id,
1768                        Some(stale_venue_order_id),
1769                        &reason,
1770                        ts,
1771                    );
1772                    return Ok(());
1773                }
1774            };
1775            let nonce = match resolve_modify_nonce(
1776                nonce_manager.next_nonce(&wallet_str, signing.subaccount_id),
1777                &emitter,
1778                strategy_id,
1779                instrument_id,
1780                client_order_id,
1781                stale_venue_order_id,
1782                clock,
1783            ) {
1784                Some(nonce) => nonce,
1785                None => return Ok(()),
1786            };
1787
1788            let payload = match order_replace_to_derive_payload(
1789                &order_for_task,
1790                &instrument,
1791                signing.subaccount_id,
1792                signing.wallet_address,
1793                &signing.signer,
1794                nonce,
1795                expiry,
1796                signing.trade_module_address,
1797                signing.domain_separator,
1798                signing.action_typehash,
1799                signing.max_fee_per_contract,
1800                Some(target_quantity.as_decimal()),
1801                target_price.map(|p| p.as_decimal()),
1802                &voi_str,
1803            ) {
1804                Ok(p) => p,
1805                Err(e) => {
1806                    let reason = format!("replace encoding failed: {e}");
1807                    log::warn!("Cannot modify order {client_order_id}: {reason}");
1808                    let ts = clock.get_time_ns();
1809                    emitter.emit_order_modify_rejected_event(
1810                        strategy_id,
1811                        instrument_id,
1812                        client_order_id,
1813                        Some(stale_venue_order_id),
1814                        &reason,
1815                        ts,
1816                    );
1817                    return Ok(());
1818                }
1819            };
1820
1821            // Mark before sending so the cancel-of-old leg is suppressed even if
1822            // it arrives before this response.
1823            dispatch_state.mark_pending_modify(client_order_id, stale_venue_order_id);
1824
1825            let outcome = ws_exec
1826                .modify_order_after_rate_limit(&payload, matching_reservation)
1827                .await;
1828
1829            if let Err(e) = &outcome
1830                && is_write_outcome_ambiguous_ws(e)
1831            {
1832                dispatch_state.clear_pending_modify(&client_order_id);
1833                log::warn!(
1834                    "Derive modify for {client_order_id} returned ambiguous WS outcome: {e}; awaiting reconciliation",
1835                );
1836                return Ok(());
1837            }
1838
1839            match outcome {
1840                Ok(DeriveReplaceOutcome::Replaced(order)) => {
1841                    let new_voi = VenueOrderId::new(order.order_id.as_str());
1842
1843                    if !dispatch_state.take_pending_modify(
1844                        &client_order_id,
1845                        stale_venue_order_id,
1846                        Some(new_voi),
1847                    ) {
1848                        log::debug!(
1849                            "Skipping private/replace response event for {client_order_id}: an incoming terminal frame already resolved the modify",
1850                        );
1851                        return Ok(());
1852                    }
1853                    log::debug!(
1854                        "Order replaced: client_order_id={client_order_id}, new venue_order_id={new_voi}",
1855                    );
1856                    let ts = clock.get_time_ns();
1857                    emitter.emit_order_updated(
1858                        &order_for_task,
1859                        new_voi,
1860                        target_quantity,
1861                        target_price,
1862                        None,
1863                        None,
1864                        ts,
1865                    );
1866                }
1867                Ok(DeriveReplaceOutcome::Canceled {
1868                    cancelled_order,
1869                    create_order_error,
1870                }) => {
1871                    if !dispatch_state.take_pending_modify(
1872                        &client_order_id,
1873                        stale_venue_order_id,
1874                        None,
1875                    ) {
1876                        log::debug!(
1877                            "Skipping partial private/replace response for {client_order_id}: an incoming terminal frame already resolved the modify",
1878                        );
1879                        return Ok(());
1880                    }
1881
1882                    log::warn!(
1883                        "Derive cancelled {client_order_id} ({}) but did not create its replacement: JSON-RPC {}: {}",
1884                        cancelled_order.order_id,
1885                        create_order_error.code,
1886                        create_order_error.message,
1887                    );
1888                    let ts = clock.get_time_ns();
1889
1890                    ensure_canceled_emitted(
1891                        &emitter,
1892                        &dispatch_state,
1893                        client_order_id,
1894                        OrderIdentity {
1895                            instrument_id,
1896                            strategy_id,
1897                            order_side: order_for_task.order_side(),
1898                            order_type: order_for_task.order_type(),
1899                        },
1900                        stale_venue_order_id,
1901                        account_id,
1902                        ts,
1903                        ts,
1904                    );
1905                    dispatch_state.forget(&client_order_id);
1906                }
1907                Err(e) => {
1908                    if !dispatch_state.take_pending_modify(
1909                        &client_order_id,
1910                        stale_venue_order_id,
1911                        None,
1912                    ) {
1913                        log::debug!(
1914                            "Skipping private/replace rejection for {client_order_id}: an incoming terminal frame already resolved the modify",
1915                        );
1916                        return Ok(());
1917                    }
1918                    let (reason, _) = ws_rejection_reason(&e);
1919                    log::debug!("Derive rejected modify for {client_order_id}: {reason}");
1920                    let ts = clock.get_time_ns();
1921                    emitter.emit_order_modify_rejected_event(
1922                        strategy_id,
1923                        instrument_id,
1924                        client_order_id,
1925                        Some(stale_venue_order_id),
1926                        &reason,
1927                        ts,
1928                    );
1929                }
1930            }
1931            Ok(())
1932        });
1933        Ok(())
1934    }
1935
1936    fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
1937        let http_client = self.http_client.clone();
1938        let subaccount_id = self.credential.subaccount_id();
1939        let emitter = self.emitter.clone();
1940        let clock = self.clock;
1941        self.spawn_task("query_account", async move {
1942            let subaccount = http_client
1943                .get_subaccount(&DeriveGetSubaccountParams::new(subaccount_id))
1944                .await?;
1945            let (balances, margins, info) = parse_derive_subaccount_to_balances(&subaccount)?;
1946            let ts_event = clock.get_time_ns();
1947            emitter.emit_account_state(balances, margins, true, ts_event, Some(info));
1948            Ok(())
1949        });
1950        Ok(())
1951    }
1952
1953    fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
1954        let Some(venue_order_id) = cmd.venue_order_id else {
1955            log::warn!(
1956                "Derive query_order requires venue_order_id (client_order_id={})",
1957                cmd.client_order_id,
1958            );
1959            return Ok(());
1960        };
1961        let http_client = self.http_client.clone();
1962        let subaccount_id = self.credential.subaccount_id();
1963        let account_id = self.core.account_id;
1964        let emitter = self.emitter.clone();
1965        let clock = self.clock;
1966        let dispatch_state = Arc::clone(&self.dispatch_state);
1967        let voi = venue_order_id.to_string();
1968
1969        self.spawn_task("query_order", async move {
1970            let order = match http_client
1971                .get_order(&DeriveGetOrderParams::new(subaccount_id, voi.as_str()))
1972                .await
1973            {
1974                Ok(o) => o,
1975                Err(e) => {
1976                    let trigger_orders = match http_client
1977                        .get_trigger_orders(&DeriveGetTriggerOrdersParams::new(subaccount_id))
1978                        .await
1979                    {
1980                        Ok(result) => result.orders,
1981                        Err(trigger_err) => {
1982                            log::warn!(
1983                                "Failed to fetch Derive order {voi}: {e}; trigger lookup also failed: {trigger_err}",
1984                            );
1985                            return Ok(());
1986                        }
1987                    };
1988
1989                    match trigger_orders
1990                        .into_iter()
1991                        .find(|o| o.order_id.as_str() == voi.as_str())
1992                    {
1993                        Some(order) => order,
1994                        None => {
1995                            log::warn!("Failed to fetch Derive order {voi}: {e}");
1996                            return Ok(());
1997                        }
1998                    }
1999                }
2000            };
2001
2002            let (price_precision, size_precision) =
2003                report_precision(&dispatch_state, order.instrument_name.as_str());
2004            let ts_init = clock.get_time_ns();
2005            let report = parse_derive_order_to_report_with_precision(
2006                &order,
2007                account_id,
2008                price_precision,
2009                size_precision,
2010                ts_init,
2011            )?;
2012            emitter.send_order_status_report(report);
2013            Ok(())
2014        });
2015        Ok(())
2016    }
2017}
2018
2019#[derive(Clone)]
2020struct DeriveReconciliationContext {
2021    http_client: DeriveHttpClient,
2022    emitter: ExecutionEventEmitter,
2023    client_id: ClientId,
2024    account_id: AccountId,
2025    subaccount_id: u64,
2026    clock: &'static AtomicTime,
2027    dispatch_state: Arc<WsDispatchState>,
2028}
2029
2030impl DeriveReconciliationContext {
2031    async fn refresh_account_state(&self) -> anyhow::Result<()> {
2032        let value = self
2033            .http_client
2034            .get_subaccount(&DeriveGetSubaccountParams::new(self.subaccount_id))
2035            .await
2036            .context("failed to fetch Derive subaccount snapshot")?;
2037        let (balances, margins, info) = parse_derive_subaccount_to_balances(&value)
2038            .context("failed to parse Derive subaccount balances")?;
2039        let ts_event = self.clock.get_time_ns();
2040        self.emitter
2041            .emit_account_state(balances, margins, true, ts_event, Some(info));
2042        Ok(())
2043    }
2044
2045    async fn recover_after_reconnect(&self) -> anyhow::Result<()> {
2046        self.refresh_account_state().await?;
2047        let mass_status = Box::pin(self.generate_mass_status(None)).await?;
2048        let order_count = mass_status.order_reports().len();
2049        let fill_count: usize = mass_status.fill_reports().values().map(Vec::len).sum();
2050        let position_count = mass_status.position_reports().len();
2051        self.emitter
2052            .send_execution_report(ExecutionReport::MassStatus(Box::new(mass_status)));
2053        log::info!(
2054            "Derive post-reconnect reconciliation submitted: orders={order_count}, fills={fill_count}, positions={position_count}",
2055        );
2056        Ok(())
2057    }
2058
2059    async fn generate_order_status_reports(
2060        &self,
2061        cmd: &GenerateOrderStatusReports,
2062        normalize_history_client_order_ids: bool,
2063    ) -> anyhow::Result<Vec<OrderStatusReport>> {
2064        let instrument_name = cmd.instrument_id.map(|id| id.symbol.as_str().to_string());
2065        let orders: Vec<DeriveOrder> = if cmd.open_only {
2066            let mut orders = self
2067                .http_client
2068                .get_open_orders(&DeriveGetOpenOrdersParams::new(self.subaccount_id))
2069                .await?
2070                .orders;
2071            orders.extend(
2072                self.http_client
2073                    .get_trigger_orders(&DeriveGetTriggerOrdersParams::new(self.subaccount_id))
2074                    .await?
2075                    .orders,
2076            );
2077            orders
2078        } else {
2079            let start_ms = cmd.start.map(|t| t.as_millis() as i64);
2080            let end_ms = cmd.end.map(|t| t.as_millis() as i64);
2081            let mut page: u32 = 1;
2082            let mut collected = Vec::new();
2083
2084            loop {
2085                let mut params = DeriveGetOrderHistoryParams::new(
2086                    self.subaccount_id,
2087                    page,
2088                    DERIVE_PRIVATE_PAGE_SIZE,
2089                )
2090                .with_window(start_ms, end_ms);
2091
2092                if let Some(name) = instrument_name.as_deref() {
2093                    params = params.with_instrument_name(name);
2094                }
2095
2096                let result = self.http_client.get_order_history(&params).await?;
2097                let total_pages = result.pagination.num_pages;
2098                collected.extend(result.orders);
2099
2100                if (page as i64) >= total_pages || total_pages == 0 {
2101                    break;
2102                }
2103                page += 1;
2104            }
2105            collected
2106        };
2107
2108        let ts_init = self.clock.get_time_ns();
2109        let start_ms = cmd.start.map(|t| t.as_millis() as i64);
2110        let end_ms = cmd.end.map(|t| t.as_millis() as i64);
2111
2112        let orders: Vec<DeriveOrder> = orders
2113            .into_iter()
2114            .filter(|order| {
2115                cmd.instrument_id.is_none_or(|instrument_id| {
2116                    InstrumentId::new(Symbol::new(order.instrument_name.as_str()), *DERIVE_VENUE)
2117                        == instrument_id
2118                }) && start_ms.is_none_or(|start| order.last_update_timestamp >= start)
2119                    && end_ms.is_none_or(|end| order.last_update_timestamp <= end)
2120            })
2121            .collect();
2122
2123        let ambiguous_client_order_ids = if normalize_history_client_order_ids {
2124            ambiguous_history_client_order_ids(&orders)
2125        } else {
2126            AHashSet::new()
2127        };
2128
2129        let mut reports = Vec::with_capacity(orders.len());
2130
2131        for order in orders {
2132            let (price_precision, size_precision) =
2133                report_precision(&self.dispatch_state, order.instrument_name.as_str());
2134            match parse_derive_order_to_report_with_precision(
2135                &order,
2136                self.account_id,
2137                price_precision,
2138                size_precision,
2139                ts_init,
2140            ) {
2141                Ok(mut report) => {
2142                    if report.client_order_id.is_some_and(|client_order_id| {
2143                        ambiguous_client_order_ids.contains(&client_order_id)
2144                    }) {
2145                        report.client_order_id = None;
2146                    }
2147                    reports.push(report);
2148                }
2149                Err(e) => log::warn!("Skipping order in status report: {e}"),
2150            }
2151        }
2152        Ok(reports)
2153    }
2154
2155    async fn generate_fill_reports(
2156        &self,
2157        cmd: GenerateFillReports,
2158    ) -> anyhow::Result<Vec<FillReport>> {
2159        let instrument_name = cmd.instrument_id.map(|id| id.symbol.as_str().to_string());
2160        let mut page: u32 = 1;
2161        let mut all_trades: Vec<DeriveTrade> = Vec::new();
2162
2163        loop {
2164            let mut params = DeriveGetTradeHistoryParams::new(
2165                self.subaccount_id,
2166                page,
2167                DERIVE_PRIVATE_PAGE_SIZE,
2168            )
2169            .with_window(
2170                cmd.start.map(|t| t.as_millis() as i64),
2171                cmd.end.map(|t| t.as_millis() as i64),
2172            );
2173
2174            if let Some(name) = instrument_name.as_deref() {
2175                params = params.with_instrument_name(name);
2176            }
2177
2178            let result = self.http_client.get_private_trade_history(&params).await?;
2179            let total_pages = result.pagination.num_pages;
2180            all_trades.extend(result.trades);
2181
2182            if (page as i64) >= total_pages || total_pages == 0 {
2183                break;
2184            }
2185            page += 1;
2186        }
2187
2188        let ts_init = self.clock.get_time_ns();
2189
2190        let venue_order_id_filter = cmd
2191            .venue_order_id
2192            .as_ref()
2193            .map(|id| id.as_str().to_string());
2194
2195        let mut reports = Vec::with_capacity(all_trades.len());
2196
2197        for trade in all_trades {
2198            if let Some(target) = venue_order_id_filter.as_deref()
2199                && trade.order_id != target
2200            {
2201                continue;
2202            }
2203
2204            let (price_precision, size_precision) =
2205                report_precision(&self.dispatch_state, trade.instrument_name.as_str());
2206            match parse_derive_trade_to_fill_report_with_precision(
2207                &trade,
2208                self.account_id,
2209                Currency::USDC(),
2210                price_precision,
2211                size_precision,
2212                ts_init,
2213            ) {
2214                Ok(Some(report)) => {
2215                    if self.dispatch_state.contains_trade(&report.trade_id) {
2216                        log::debug!(
2217                            "Skipping duplicate Derive fill (trade_id={}) in generate_fill_reports",
2218                            report.trade_id,
2219                        );
2220                        continue;
2221                    }
2222                    reports.push(report);
2223                }
2224                Ok(None) => {}
2225                Err(e) => log::warn!("Skipping trade in fill report: {e}"),
2226            }
2227        }
2228        Ok(reports)
2229    }
2230
2231    async fn generate_position_status_snapshot(
2232        &self,
2233        cmd: &GeneratePositionStatusReports,
2234    ) -> anyhow::Result<PositionStatusSnapshot> {
2235        let positions = self
2236            .http_client
2237            .get_positions(&DeriveGetPositionsParams::new(self.subaccount_id))
2238            .await?
2239            .positions;
2240        let ts_init = self.clock.get_time_ns();
2241        let mut reports = Vec::with_capacity(positions.len());
2242        let mut instruments = AHashSet::with_capacity(positions.len());
2243
2244        for position in positions {
2245            let instrument_id = format_instrument_id(position.instrument_name.as_str());
2246            if let Some(target) = cmd.instrument_id
2247                && instrument_id != target
2248            {
2249                continue;
2250            }
2251
2252            instruments.insert(instrument_id);
2253
2254            let (_, size_precision) =
2255                report_precision(&self.dispatch_state, position.instrument_name.as_str());
2256            match parse_derive_position_to_report_with_precision(
2257                &position,
2258                self.account_id,
2259                size_precision,
2260                ts_init,
2261            ) {
2262                Ok(report) => reports.push(report),
2263                Err(e) => log::warn!("Skipping position in status report: {e}"),
2264            }
2265        }
2266
2267        Ok(PositionStatusSnapshot {
2268            reports,
2269            instruments,
2270        })
2271    }
2272
2273    async fn generate_mass_status(
2274        &self,
2275        lookback_mins: Option<u64>,
2276    ) -> anyhow::Result<ExecutionMassStatus> {
2277        log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
2278
2279        let ts_now = self.clock.get_time_ns();
2280        let start = lookback_mins.map(|mins| {
2281            let lookback_ns = mins.saturating_mul(60).saturating_mul(1_000_000_000);
2282            UnixNanos::from(ts_now.as_u64().saturating_sub(lookback_ns))
2283        });
2284        let open_order_cmd = GenerateOrderStatusReports::new(
2285            UUID4::new(),
2286            ts_now,
2287            true,
2288            None,
2289            None,
2290            None,
2291            None,
2292            None,
2293        );
2294        let history_order_cmd = GenerateOrderStatusReports::new(
2295            UUID4::new(),
2296            ts_now,
2297            false,
2298            None,
2299            start,
2300            None,
2301            None,
2302            None,
2303        );
2304        let fill_cmd =
2305            GenerateFillReports::new(UUID4::new(), ts_now, None, None, start, None, None, None);
2306        let position_cmd =
2307            GeneratePositionStatusReports::new(UUID4::new(), ts_now, None, None, None, None, None);
2308
2309        let (history_order_reports, open_order_reports, mut fill_reports, position_snapshot) = tokio::try_join!(
2310            self.generate_order_status_reports(&history_order_cmd, true),
2311            self.generate_order_status_reports(&open_order_cmd, false),
2312            self.generate_fill_reports(fill_cmd),
2313            self.generate_position_status_snapshot(&position_cmd),
2314        )?;
2315        let detached_history_order_ids: AHashSet<VenueOrderId> = history_order_reports
2316            .iter()
2317            .filter(|report| report.client_order_id.is_none())
2318            .map(|report| report.venue_order_id)
2319            .collect();
2320
2321        for report in &mut fill_reports {
2322            if detached_history_order_ids.contains(&report.venue_order_id) {
2323                report.client_order_id = None;
2324            }
2325        }
2326
2327        log::info!(
2328            "Received {} historical OrderStatusReports",
2329            history_order_reports.len()
2330        );
2331        log::info!(
2332            "Received {} open OrderStatusReports",
2333            open_order_reports.len()
2334        );
2335        log::info!("Received {} FillReports", fill_reports.len());
2336        log::info!(
2337            "Received {} PositionReports",
2338            position_snapshot.reports.len()
2339        );
2340
2341        let mut touched_instruments = AHashSet::new();
2342
2343        for report in history_order_reports
2344            .iter()
2345            .chain(open_order_reports.iter())
2346        {
2347            touched_instruments.insert(report.instrument_id);
2348        }
2349
2350        for report in &fill_reports {
2351            touched_instruments.insert(report.instrument_id);
2352        }
2353
2354        let PositionStatusSnapshot {
2355            reports: position_reports,
2356            instruments: position_instruments,
2357        } = position_snapshot;
2358        let mut mass_status =
2359            ExecutionMassStatus::new(self.client_id, self.account_id, *DERIVE_VENUE, ts_now, None);
2360        mass_status.add_order_reports(history_order_reports);
2361        mass_status.add_order_reports(open_order_reports);
2362        mass_status.add_fill_reports(fill_reports);
2363        mass_status.add_position_reports(position_reports);
2364
2365        add_missing_flat_position_reports(
2366            &mut mass_status,
2367            self.account_id,
2368            touched_instruments,
2369            &position_instruments,
2370            ts_now,
2371        );
2372
2373        Ok(mass_status)
2374    }
2375}
2376
2377fn ambiguous_history_client_order_ids(orders: &[DeriveOrder]) -> AHashSet<ClientOrderId> {
2378    let mut orders_by_label: AHashMap<Ustr, AHashMap<&str, Option<&str>>> = AHashMap::new();
2379
2380    for order in orders {
2381        if order.label.is_empty() {
2382            continue;
2383        }
2384        orders_by_label
2385            .entry(order.label)
2386            .or_default()
2387            .insert(order.order_id.as_str(), order.replaced_order_id.as_deref());
2388    }
2389
2390    let mut ambiguous_client_order_ids = AHashSet::new();
2391
2392    for (label, orders_by_id) in orders_by_label {
2393        if orders_by_id.len() < 2 {
2394            continue;
2395        }
2396
2397        let predecessors: AHashMap<&str, &str> = orders_by_id
2398            .iter()
2399            .filter_map(|(order_id, replaced_order_id)| {
2400                let replaced_order_id = (*replaced_order_id)?;
2401                orders_by_id
2402                    .contains_key(replaced_order_id)
2403                    .then_some((*order_id, replaced_order_id))
2404            })
2405            .collect();
2406        let predecessor_ids: AHashSet<&str> = predecessors.values().copied().collect();
2407        let heads: Vec<&str> = orders_by_id
2408            .keys()
2409            .copied()
2410            .filter(|order_id| !predecessor_ids.contains(order_id))
2411            .collect();
2412
2413        // One client order may own several venue IDs only when they form one linear replace chain
2414        let is_linear_chain = predecessors.len() + 1 == orders_by_id.len()
2415            && predecessor_ids.len() == predecessors.len()
2416            && heads.len() == 1
2417            && {
2418                let mut visited = AHashSet::new();
2419                let mut current = Some(heads[0]);
2420                while let Some(order_id) = current {
2421                    if !visited.insert(order_id) {
2422                        break;
2423                    }
2424                    current = predecessors.get(order_id).copied();
2425                }
2426                visited.len() == orders_by_id.len()
2427            };
2428
2429        if !is_linear_chain {
2430            ambiguous_client_order_ids.insert(ClientOrderId::new(label.as_str()));
2431        }
2432    }
2433
2434    ambiguous_client_order_ids
2435}
2436
2437struct PositionStatusSnapshot {
2438    reports: Vec<PositionStatusReport>,
2439    instruments: AHashSet<InstrumentId>,
2440}
2441
2442// Reason text and post-only classification for a definitive WS write failure.
2443// Non-JSON-RPC errors carry no venue code and are never post-only crossings.
2444fn ws_rejection_reason(error: &DeriveWsError) -> (String, bool) {
2445    match error {
2446        DeriveWsError::JsonRpc { code, message, .. } => (
2447            format!("JSON-RPC {code}: {message}"),
2448            derive_rejection_due_post_only(Some(*code), message),
2449        ),
2450        other => (other.to_string(), false),
2451    }
2452}
2453
2454fn add_missing_flat_position_reports(
2455    mass_status: &mut ExecutionMassStatus,
2456    account_id: AccountId,
2457    touched_instruments: AHashSet<InstrumentId>,
2458    position_instruments: &AHashSet<InstrumentId>,
2459    ts_init: UnixNanos,
2460) {
2461    let mut flat_reports = Vec::new();
2462
2463    for instrument_id in touched_instruments {
2464        if position_instruments.contains(&instrument_id) {
2465            continue;
2466        }
2467
2468        flat_reports.push(PositionStatusReport::new(
2469            account_id,
2470            instrument_id,
2471            PositionSide::Flat,
2472            Quantity::from("0"),
2473            ts_init,
2474            ts_init,
2475            Some(UUID4::new()),
2476            None,
2477            None,
2478        ));
2479    }
2480
2481    if !flat_reports.is_empty() {
2482        log::info!(
2483            "Added {} flat PositionReports for Derive instruments absent from current positions",
2484            flat_reports.len()
2485        );
2486        mass_status.add_position_reports(flat_reports);
2487    }
2488}
2489
2490fn report_precision(
2491    dispatch_state: &WsDispatchState,
2492    instrument_name: &str,
2493) -> (Option<u8>, Option<u8>) {
2494    let instrument_id = format_instrument_id(instrument_name);
2495    dispatch_state
2496        .instrument_precision(&instrument_id)
2497        .map_or((None, None), |(price, size)| (Some(price), Some(size)))
2498}
2499
2500fn handle_ws_message(
2501    message: DeriveWsMessage,
2502    emitter: &ExecutionEventEmitter,
2503    account_id: AccountId,
2504    clock: &'static AtomicTime,
2505    dispatch_state: &WsDispatchState,
2506) {
2507    let payload = match message {
2508        DeriveWsMessage::Subscription(payload) => payload,
2509        DeriveWsMessage::Authenticated
2510        | DeriveWsMessage::Reconnected
2511        | DeriveWsMessage::SessionRecoveryFailed(_) => return,
2512    };
2513
2514    let is_orders_channel = payload.channel.as_str().ends_with(".orders");
2515    let is_trades_channel = payload.channel.as_str().ends_with(".trades");
2516
2517    if is_orders_channel {
2518        let data = match serde_json::from_str::<DeriveOrdersSubscriptionData>(payload.data.get()) {
2519            Ok(data) => data,
2520            Err(e) => {
2521                log::warn!(
2522                    "Failed to decode Derive orders frame on channel {}: {e}",
2523                    payload.channel,
2524                );
2525                return;
2526            }
2527        };
2528        dispatch_orders_payload(data, emitter, account_id, clock, dispatch_state);
2529    } else if is_trades_channel {
2530        let data = match serde_json::from_str::<DeriveTradesSubscriptionData>(payload.data.get()) {
2531            Ok(data) => data,
2532            Err(e) => {
2533                log::warn!(
2534                    "Failed to decode Derive trades frame on channel {}: {e}",
2535                    payload.channel,
2536                );
2537                return;
2538            }
2539        };
2540        dispatch_trades_payload(data, emitter, account_id, clock, dispatch_state);
2541    }
2542}
2543
2544/// Dispatches a parsed `{subaccount_id}.orders` payload to the execution event
2545/// emitter.
2546///
2547/// Emits tracked order events when an order's client order id resolves to a
2548/// registered identity in `dispatch_state`, and forwards a raw status report
2549/// otherwise.
2550pub fn dispatch_orders_payload(
2551    data: DeriveOrdersSubscriptionData,
2552    emitter: &ExecutionEventEmitter,
2553    account_id: AccountId,
2554    clock: &'static AtomicTime,
2555    dispatch_state: &WsDispatchState,
2556) {
2557    let ts_init = clock.get_time_ns();
2558
2559    for order in data.orders {
2560        let (price_precision, size_precision) =
2561            report_precision(dispatch_state, order.instrument_name.as_str());
2562        let report = match parse_derive_order_to_report_with_precision(
2563            &order,
2564            account_id,
2565            price_precision,
2566            size_precision,
2567            ts_init,
2568        ) {
2569            Ok(report) => report,
2570            Err(e) => {
2571                log::warn!("Failed to parse Derive order WS update: {e}");
2572                continue;
2573            }
2574        };
2575
2576        let identity = tracked_order_identity(report.client_order_id, dispatch_state);
2577
2578        match identity {
2579            Some((client_order_id, identity)) => emit_tracked_order_event(
2580                emitter,
2581                dispatch_state,
2582                client_order_id,
2583                identity,
2584                &report,
2585                account_id,
2586                ts_init,
2587            ),
2588            None => emitter.send_order_status_report(report),
2589        }
2590    }
2591}
2592
2593/// Dispatches a parsed `{subaccount_id}.trades` payload to the execution event
2594/// emitter.
2595///
2596/// Deduplicates by trade id, then emits a tracked fill when the trade's client
2597/// order id resolves to a registered identity in `dispatch_state`, and forwards
2598/// a raw fill report otherwise.
2599pub fn dispatch_trades_payload(
2600    data: DeriveTradesSubscriptionData,
2601    emitter: &ExecutionEventEmitter,
2602    account_id: AccountId,
2603    clock: &'static AtomicTime,
2604    dispatch_state: &WsDispatchState,
2605) {
2606    let fee_currency = Currency::USDC();
2607    let ts_init = clock.get_time_ns();
2608
2609    for trade in data.trades {
2610        let (price_precision, size_precision) =
2611            report_precision(dispatch_state, trade.instrument_name.as_str());
2612        match parse_derive_trade_to_fill_report_with_precision(
2613            &trade,
2614            account_id,
2615            fee_currency,
2616            price_precision,
2617            size_precision,
2618            ts_init,
2619        ) {
2620            Ok(Some(report)) => {
2621                if dispatch_state.check_and_insert_trade(report.trade_id) {
2622                    log::debug!(
2623                        "Skipping duplicate Derive fill (trade_id={}) on WS dispatch",
2624                        report.trade_id,
2625                    );
2626                    continue;
2627                }
2628
2629                let identity = tracked_order_identity(report.client_order_id, dispatch_state);
2630
2631                match identity {
2632                    Some((client_order_id, identity)) => emit_tracked_fill(
2633                        emitter,
2634                        dispatch_state,
2635                        client_order_id,
2636                        identity,
2637                        &report,
2638                        account_id,
2639                        ts_init,
2640                    ),
2641                    None => emitter.send_fill_report(report),
2642                }
2643            }
2644            Ok(None) => {}
2645            Err(e) => log::warn!("Failed to parse Derive trade WS update: {e}"),
2646        }
2647    }
2648}
2649
2650fn tracked_order_identity(
2651    client_order_id: Option<ClientOrderId>,
2652    dispatch_state: &WsDispatchState,
2653) -> Option<(ClientOrderId, OrderIdentity)> {
2654    client_order_id.and_then(|cid| {
2655        dispatch_state
2656            .identity(&cid)
2657            .map(|identity| (cid, identity))
2658    })
2659}
2660
2661/// Synthesizes and emits `OrderAccepted` when one has not yet been emitted
2662/// for the order. Used to guarantee the `Submitted -> Accepted -> ...`
2663/// lifecycle when a fill or terminal event arrives before (or instead of)
2664/// the venue's `Open` notice.
2665#[expect(clippy::too_many_arguments)]
2666fn ensure_accepted_emitted(
2667    emitter: &ExecutionEventEmitter,
2668    dispatch_state: &WsDispatchState,
2669    client_order_id: ClientOrderId,
2670    identity: OrderIdentity,
2671    venue_order_id: VenueOrderId,
2672    account_id: AccountId,
2673    ts_event: UnixNanos,
2674    ts_init: UnixNanos,
2675) {
2676    if dispatch_state.mark_accepted(client_order_id) {
2677        return;
2678    }
2679    let accepted = OrderAccepted::new(
2680        emitter.trader_id(),
2681        identity.strategy_id,
2682        identity.instrument_id,
2683        client_order_id,
2684        venue_order_id,
2685        account_id,
2686        UUID4::new(),
2687        ts_event,
2688        ts_init,
2689        false,
2690    );
2691    emitter.send_order_event(OrderEventAny::Accepted(accepted));
2692}
2693
2694#[expect(clippy::too_many_arguments)]
2695fn ensure_canceled_emitted(
2696    emitter: &ExecutionEventEmitter,
2697    dispatch_state: &WsDispatchState,
2698    client_order_id: ClientOrderId,
2699    identity: OrderIdentity,
2700    venue_order_id: VenueOrderId,
2701    account_id: AccountId,
2702    ts_event: UnixNanos,
2703    ts_init: UnixNanos,
2704) {
2705    if dispatch_state.mark_canceled(client_order_id) {
2706        return;
2707    }
2708    let canceled = OrderCanceled::new(
2709        emitter.trader_id(),
2710        identity.strategy_id,
2711        identity.instrument_id,
2712        client_order_id,
2713        UUID4::new(),
2714        ts_event,
2715        ts_init,
2716        false,
2717        Some(venue_order_id),
2718        Some(account_id),
2719    );
2720    emitter.send_order_event(OrderEventAny::Canceled(canceled));
2721}
2722
2723fn emit_tracked_order_event(
2724    emitter: &ExecutionEventEmitter,
2725    dispatch_state: &WsDispatchState,
2726    client_order_id: ClientOrderId,
2727    identity: OrderIdentity,
2728    report: &OrderStatusReport,
2729    account_id: AccountId,
2730    ts_init: UnixNanos,
2731) {
2732    let venue_order_id = report.venue_order_id;
2733    let ts_accepted = report.ts_accepted;
2734    let ts_event = report.ts_last;
2735
2736    // A `private/replace` cancels the old order and opens a new one under the
2737    // same label; suppress events for the superseded old venue order id so they
2738    // don't terminate the order that `modify_order` rebinds via `OrderUpdated`.
2739    // `pending_modify` covers the in-flight window; the bound-id check covers
2740    // after the rebind.
2741    if dispatch_state.pending_modify(&client_order_id) == Some(venue_order_id) {
2742        log::debug!(
2743            "Skipping cancel-replace leg for {client_order_id}: stale venue_order_id={venue_order_id}",
2744        );
2745        return;
2746    }
2747
2748    if let Some(bound) = dispatch_state.bound_venue_order_id(&client_order_id)
2749        && bound != venue_order_id
2750    {
2751        let terminal = matches!(
2752            report.order_status,
2753            OrderStatus::Canceled | OrderStatus::Expired | OrderStatus::Rejected
2754        );
2755
2756        if dispatch_state.bind_incoming_modify(client_order_id, venue_order_id, terminal) {
2757            log::debug!(
2758                "Bound incoming replacement for {client_order_id}: venue_order_id={venue_order_id}",
2759            );
2760        } else {
2761            log::debug!(
2762                "Skipping stale {:?} for {client_order_id}: venue_order_id={venue_order_id} superseded by {bound}",
2763                report.order_status,
2764            );
2765            return;
2766        }
2767    }
2768
2769    match report.order_status {
2770        OrderStatus::Accepted | OrderStatus::PartiallyFilled => {
2771            if dispatch_state.contains_filled(&client_order_id) {
2772                log::debug!("Skipping stale Accepted for {client_order_id} (already filled)",);
2773                return;
2774            }
2775            dispatch_state.record_venue_order_id(client_order_id, venue_order_id);
2776            ensure_accepted_emitted(
2777                emitter,
2778                dispatch_state,
2779                client_order_id,
2780                identity,
2781                venue_order_id,
2782                account_id,
2783                ts_accepted,
2784                ts_init,
2785            );
2786        }
2787        OrderStatus::Filled => {
2788            dispatch_state.record_venue_order_id(client_order_id, venue_order_id);
2789            ensure_accepted_emitted(
2790                emitter,
2791                dispatch_state,
2792                client_order_id,
2793                identity,
2794                venue_order_id,
2795                account_id,
2796                ts_accepted,
2797                ts_init,
2798            );
2799            // Mark the order terminal so replayed Accepted frames are
2800            // suppressed, but keep its identity alive: the matching
2801            // `.trades` frame may arrive after this `.orders` Filled
2802            // notice and still needs the tracked path to emit a proper
2803            // `OrderFilled`. Identity is retired by Canceled/Expired/
2804            // Rejected paths; full-fill leaks are bounded by submission
2805            // throughput.
2806            dispatch_state.mark_filled(client_order_id);
2807        }
2808        OrderStatus::Canceled => {
2809            ensure_accepted_emitted(
2810                emitter,
2811                dispatch_state,
2812                client_order_id,
2813                identity,
2814                venue_order_id,
2815                account_id,
2816                ts_accepted,
2817                ts_init,
2818            );
2819            ensure_canceled_emitted(
2820                emitter,
2821                dispatch_state,
2822                client_order_id,
2823                identity,
2824                venue_order_id,
2825                account_id,
2826                ts_event,
2827                ts_init,
2828            );
2829            dispatch_state.forget(&client_order_id);
2830        }
2831        OrderStatus::Expired => {
2832            ensure_accepted_emitted(
2833                emitter,
2834                dispatch_state,
2835                client_order_id,
2836                identity,
2837                venue_order_id,
2838                account_id,
2839                ts_accepted,
2840                ts_init,
2841            );
2842            let expired = OrderExpired::new(
2843                emitter.trader_id(),
2844                identity.strategy_id,
2845                identity.instrument_id,
2846                client_order_id,
2847                UUID4::new(),
2848                ts_event,
2849                ts_init,
2850                false,
2851                Some(venue_order_id),
2852                Some(account_id),
2853            );
2854            emitter.send_order_event(OrderEventAny::Expired(expired));
2855            dispatch_state.forget(&client_order_id);
2856        }
2857        OrderStatus::Rejected => {
2858            let reason = report
2859                .cancel_reason
2860                .as_deref()
2861                .unwrap_or("Order rejected by Derive");
2862            let due_post_only = derive_rejection_due_post_only(None, reason);
2863            let rejected = OrderRejected::new(
2864                emitter.trader_id(),
2865                identity.strategy_id,
2866                identity.instrument_id,
2867                client_order_id,
2868                account_id,
2869                Ustr::from(reason),
2870                UUID4::new(),
2871                ts_event,
2872                ts_init,
2873                false,
2874                due_post_only,
2875            );
2876            emitter.send_order_event(OrderEventAny::Rejected(rejected));
2877            dispatch_state.forget(&client_order_id);
2878        }
2879        other => {
2880            log::debug!(
2881                "Unhandled tracked order status {other:?} for {client_order_id}, sending as report",
2882            );
2883            emitter.send_order_status_report(report.clone());
2884        }
2885    }
2886}
2887
2888fn emit_tracked_fill(
2889    emitter: &ExecutionEventEmitter,
2890    dispatch_state: &WsDispatchState,
2891    client_order_id: ClientOrderId,
2892    identity: OrderIdentity,
2893    report: &FillReport,
2894    account_id: AccountId,
2895    ts_init: UnixNanos,
2896) {
2897    ensure_accepted_emitted(
2898        emitter,
2899        dispatch_state,
2900        client_order_id,
2901        identity,
2902        report.venue_order_id,
2903        account_id,
2904        report.ts_event,
2905        ts_init,
2906    );
2907
2908    let filled = OrderFilled::new(
2909        emitter.trader_id(),
2910        identity.strategy_id,
2911        identity.instrument_id,
2912        client_order_id,
2913        report.venue_order_id,
2914        account_id,
2915        report.trade_id,
2916        identity.order_side,
2917        identity.order_type,
2918        report.last_qty,
2919        report.last_px,
2920        report.commission.currency,
2921        report.liquidity_side,
2922        UUID4::new(),
2923        report.ts_event,
2924        ts_init,
2925        false,
2926        report.venue_position_id,
2927        Some(report.commission),
2928        None,
2929    );
2930    emitter.send_order_event(OrderEventAny::Filled(filled));
2931}
2932
2933/// Derives the worst-acceptable limit price for a market order from the
2934/// top-of-book quote and a slippage bound in basis points, rounded to the
2935/// instrument's `tick_size`.
2936///
2937/// Buys lift the ask by `slippage_bps` then round up to the next tick; sells
2938/// drop the bid by the same and round down. The result is the signed
2939/// `limit_price` slot in the EIP-712 trade module data; the venue uses it
2940/// as a worst-case bound while the order sweeps. A non-positive sell bound
2941/// is rejected (`None`) so the caller can deny the order rather than sign
2942/// an invalid zero limit.
2943fn market_order_limit_price(
2944    quote: &QuoteTick,
2945    side: OrderSide,
2946    slippage_bps: u32,
2947    tick_size: Decimal,
2948) -> Option<Decimal> {
2949    let bps = Decimal::from(slippage_bps);
2950    let scale = Decimal::from(10_000_u32);
2951    let one = Decimal::ONE;
2952    let raw = match side {
2953        OrderSide::Buy => quote.ask_price.as_decimal() * (one + bps / scale),
2954        OrderSide::Sell => quote.bid_price.as_decimal() * (one - bps / scale),
2955    };
2956    let rounded = round_to_tick(raw, tick_size, side);
2957    if rounded <= Decimal::ZERO {
2958        return None;
2959    }
2960    Some(rounded)
2961}
2962
2963fn trigger_market_limit_price(
2964    trigger_price: Decimal,
2965    side: OrderSide,
2966    slippage_bps: u32,
2967    tick_size: Decimal,
2968) -> Option<Decimal> {
2969    let bps = Decimal::from(slippage_bps);
2970    let scale = Decimal::from(10_000_u32);
2971    let one = Decimal::ONE;
2972    let raw = match side {
2973        OrderSide::Buy => trigger_price * (one + bps / scale),
2974        OrderSide::Sell => trigger_price * (one - bps / scale),
2975    };
2976    let rounded = round_to_tick(raw, tick_size, side);
2977    if rounded <= Decimal::ZERO {
2978        return None;
2979    }
2980    Some(rounded)
2981}
2982
2983fn is_derive_trigger_order_type(order_type: OrderType) -> bool {
2984    matches!(
2985        order_type,
2986        OrderType::StopMarket
2987            | OrderType::StopLimit
2988            | OrderType::MarketIfTouched
2989            | OrderType::LimitIfTouched
2990    )
2991}
2992
2993fn trigger_order_signature_expiry(clock: &'static AtomicTime) -> i64 {
2994    let now_secs = (clock.get_time_ns().as_u64() / 1_000_000_000) as i64;
2995    now_secs + TRIGGER_ORDER_SIGNATURE_TTL.as_secs() as i64
2996}
2997
2998fn resolve_submit_nonce(
2999    nonce: Result<u64, NonceError>,
3000    emitter: &ExecutionEventEmitter,
3001    dispatch_state: &WsDispatchState,
3002    order: &OrderAny,
3003    clock: &'static AtomicTime,
3004) -> Option<u64> {
3005    match nonce {
3006        Ok(nonce) => Some(nonce),
3007        Err(e) => {
3008            let reason = format!("nonce allocation failed: {e}");
3009            log::warn!("Cannot submit order {}: {reason}", order.client_order_id());
3010            dispatch_state.forget(&order.client_order_id());
3011            emitter.emit_order_rejected(order, &reason, clock.get_time_ns(), false);
3012            None
3013        }
3014    }
3015}
3016
3017fn resolve_modify_nonce(
3018    nonce: Result<u64, NonceError>,
3019    emitter: &ExecutionEventEmitter,
3020    strategy_id: StrategyId,
3021    instrument_id: InstrumentId,
3022    client_order_id: ClientOrderId,
3023    venue_order_id: VenueOrderId,
3024    clock: &'static AtomicTime,
3025) -> Option<u64> {
3026    match nonce {
3027        Ok(nonce) => Some(nonce),
3028        Err(e) => {
3029            let reason = format!("nonce allocation failed: {e}");
3030            log::warn!("Cannot modify order {client_order_id}: {reason}");
3031            emitter.emit_order_modify_rejected_event(
3032                strategy_id,
3033                instrument_id,
3034                client_order_id,
3035                Some(venue_order_id),
3036                &reason,
3037                clock.get_time_ns(),
3038            );
3039            None
3040        }
3041    }
3042}
3043
3044fn normal_order_signature_expiry(
3045    clock: &'static AtomicTime,
3046    signature_expiry_secs: u64,
3047) -> anyhow::Result<i64> {
3048    let min_ttl_secs = MIN_SIGNATURE_TTL.as_secs();
3049    if signature_expiry_secs <= min_ttl_secs {
3050        anyhow::bail!(
3051            "signature_expiry_secs {signature_expiry_secs}s must be greater than the Derive minimum {min_ttl_secs}s"
3052        );
3053    }
3054
3055    let now_secs_u64 = clock.get_time_ns().as_u64() / 1_000_000_000;
3056    let now_secs = i64::try_from(now_secs_u64).with_context(|| {
3057        format!("current UNIX time {now_secs_u64}s cannot fit in Derive signature_expiry_sec")
3058    })?;
3059    let ttl_secs = i64::try_from(signature_expiry_secs).with_context(|| {
3060        format!(
3061            "signature_expiry_secs {signature_expiry_secs}s cannot fit in Derive signature_expiry_sec"
3062        )
3063    })?;
3064
3065    now_secs.checked_add(ttl_secs).ok_or_else(|| {
3066        anyhow::anyhow!(
3067            "signature expiry overflows Derive signature_expiry_sec: now {now_secs}s plus TTL {ttl_secs}s"
3068        )
3069    })
3070}
3071
3072async fn refresh_market_order_quote(
3073    http_client: &DeriveHttpClient,
3074    venue_symbol: &str,
3075    instrument: &DeriveInstrument,
3076    clock: &'static AtomicTime,
3077) -> anyhow::Result<QuoteTick> {
3078    let ticker = http_client.get_ticker(venue_symbol).await?;
3079    let price_precision = Price::from_decimal(instrument.tick_size)
3080        .with_context(|| format!("invalid Derive tick_size for {venue_symbol}"))?
3081        .precision;
3082    let size_precision = Quantity::from_decimal(instrument.amount_step)
3083        .with_context(|| format!("invalid Derive amount_step for {venue_symbol}"))?
3084        .precision;
3085
3086    parse_ticker_quote_from_rest(
3087        &ticker,
3088        price_precision,
3089        size_precision,
3090        clock.get_time_ns(),
3091    )
3092}
3093
3094/// Rounds `value` to the nearest multiple of `tick_size`. Buys round up so
3095/// the signed bound remains acceptable to the venue; sells round down so the
3096/// caller does not accidentally tighten the floor. A non-positive `tick_size`
3097/// is treated as a no-op.
3098fn round_to_tick(value: Decimal, tick_size: Decimal, side: OrderSide) -> Decimal {
3099    if tick_size <= Decimal::ZERO {
3100        return value;
3101    }
3102    let ratio = value / tick_size;
3103    let ticks = match side {
3104        OrderSide::Buy => ratio.ceil(),
3105        OrderSide::Sell => ratio.floor(),
3106    };
3107    ticks * tick_size
3108}
3109
3110async fn cached_or_fetch_instrument(
3111    http_client: &DeriveHttpClient,
3112    instruments: &Arc<AtomicMap<InstrumentId, DeriveInstrument>>,
3113    instrument_id: &InstrumentId,
3114    venue_symbol: &str,
3115) -> anyhow::Result<DeriveInstrument> {
3116    if let Some(cached) = instruments.get_cloned(instrument_id) {
3117        return Ok(cached);
3118    }
3119    let instrument = http_client
3120        .get_instrument(venue_symbol)
3121        .await
3122        .with_context(|| format!("failed to fetch instrument {venue_symbol}"))?;
3123    instruments.insert(*instrument_id, instrument.clone());
3124    Ok(instrument)
3125}
3126
3127#[cfg(test)]
3128mod tests {
3129    use std::{cell::RefCell, rc::Rc};
3130
3131    use nautilus_common::{
3132        cache::Cache,
3133        messages::{ExecutionEvent, ExecutionReport},
3134    };
3135    use nautilus_core::UnixNanos;
3136    use nautilus_live::ExecutionClientCore;
3137    use nautilus_model::{
3138        data::QuoteTick,
3139        enums::{AccountType, OmsType, TimeInForce},
3140        identifiers::{AccountId, ClientId, InstrumentId, StrategyId, TraderId},
3141        orders::OrderTestBuilder,
3142        types::{Price, Quantity},
3143    };
3144    use rstest::rstest;
3145    use rust_decimal_macros::dec;
3146
3147    use super::*;
3148    use crate::common::{
3149        consts::DERIVE,
3150        enums::{DeriveEnvironment, DeriveOrderStatus, DeriveOrderType},
3151        parse::parse_derive_instrument_any,
3152    };
3153
3154    const TEST_WALLET: &str = "0x0000000000000000000000000000000000001234";
3155    const TEST_SESSION_KEY: &str =
3156        "0x2ae8be44db8a590d20bffbe3b6872df9b569147d3bf6801a35a28281a4816bbd";
3157    const TEST_SUBACCOUNT: u64 = 30769;
3158
3159    fn test_core() -> ExecutionClientCore {
3160        let cache = Rc::new(RefCell::new(Cache::default()));
3161        ExecutionClientCore::new(
3162            TraderId::from("TRADER-001"),
3163            ClientId::from(DERIVE),
3164            *DERIVE_VENUE,
3165            OmsType::Netting,
3166            AccountId::from("DERIVE-001"),
3167            AccountType::Margin,
3168            None,
3169            cache,
3170        )
3171    }
3172
3173    fn test_config() -> DeriveExecutionClientConfig {
3174        DeriveExecutionClientConfig {
3175            wallet_address: Some(TEST_WALLET.to_string()),
3176            session_key: Some(TEST_SESSION_KEY.to_string()),
3177            subaccount_id: Some(TEST_SUBACCOUNT),
3178            environment: DeriveEnvironment::Testnet,
3179            domain_separator: Some(
3180                "0x2222222222222222222222222222222222222222222222222222222222222222".to_string(),
3181            ),
3182            action_typehash: Some(
3183                "0x1111111111111111111111111111111111111111111111111111111111111111".to_string(),
3184            ),
3185            trade_module_address: Some("0x000000000000000000000000000000000000bbbb".to_string()),
3186            max_fee_per_contract: Some(dec!(1000)),
3187            ..DeriveExecutionClientConfig::default()
3188        }
3189    }
3190
3191    #[rstest]
3192    fn test_market_order_limit_price_buy_lifts_ask_and_rounds_up_to_tick() {
3193        let quote = QuoteTick::new(
3194            InstrumentId::from("ETH-PERP.DERIVE"),
3195            Price::from("3500.00"),
3196            Price::from("3501.00"),
3197            Quantity::from("1.000"),
3198            Quantity::from("1.000"),
3199            UnixNanos::from(0),
3200            UnixNanos::from(0),
3201        );
3202        // 50 bps; raw = 3501 * 1.005 = 3518.505; tick 0.01 rounds up to 3518.51.
3203        let price = market_order_limit_price(&quote, OrderSide::Buy, 50, dec!(0.01)).unwrap();
3204        assert_eq!(price, dec!(3518.51));
3205    }
3206
3207    #[rstest]
3208    fn test_market_order_limit_price_sell_drops_bid_rounds_down_and_denies_non_positive() {
3209        let quote = QuoteTick::new(
3210            InstrumentId::from("ETH-PERP.DERIVE"),
3211            Price::from("3500.00"),
3212            Price::from("3501.00"),
3213            Quantity::from("1.000"),
3214            Quantity::from("1.000"),
3215            UnixNanos::from(0),
3216            UnixNanos::from(0),
3217        );
3218        // 50 bps; raw = 3500 * 0.995 = 3482.5; tick 0.01 stays at 3482.5.
3219        let price = market_order_limit_price(&quote, OrderSide::Sell, 50, dec!(0.01)).unwrap();
3220        assert_eq!(price, dec!(3482.5));
3221
3222        // 20_000 bps = 200% slippage drives the rounded bound below zero; deny.
3223        let zero = market_order_limit_price(&quote, OrderSide::Sell, 20_000, dec!(0.01));
3224        assert!(zero.is_none());
3225    }
3226
3227    #[rstest]
3228    fn test_trigger_market_limit_price_uses_trigger_price_bound() {
3229        let buy = trigger_market_limit_price(dec!(3600), OrderSide::Buy, 50, dec!(0.01)).unwrap();
3230        let sell = trigger_market_limit_price(dec!(3600), OrderSide::Sell, 50, dec!(0.01)).unwrap();
3231        let zero = trigger_market_limit_price(dec!(1), OrderSide::Sell, 20_000, dec!(0.01));
3232
3233        assert_eq!(buy, dec!(3618));
3234        assert_eq!(sell, dec!(3582));
3235        assert!(zero.is_none());
3236    }
3237
3238    #[rstest]
3239    fn test_normal_order_signature_expiry_accepts_ttl_above_minimum() {
3240        let clock = get_atomic_clock_realtime();
3241        let start_secs = (clock.get_time_ns().as_u64() / 1_000_000_000) as i64;
3242        let ttl_secs = MIN_SIGNATURE_TTL.as_secs() + 1;
3243
3244        let expiry = normal_order_signature_expiry(clock, ttl_secs).expect("expiry is valid");
3245
3246        assert!(expiry >= start_secs + ttl_secs as i64);
3247    }
3248
3249    #[rstest]
3250    #[case(MIN_SIGNATURE_TTL.as_secs(), "must be greater than the Derive minimum")]
3251    #[case(MIN_SIGNATURE_TTL.as_secs() - 1, "must be greater than the Derive minimum")]
3252    fn test_normal_order_signature_expiry_rejects_minimum_or_lower_ttl(
3253        #[case] ttl_secs: u64,
3254        #[case] reason_fragment: &str,
3255    ) {
3256        let clock = get_atomic_clock_realtime();
3257
3258        let err = normal_order_signature_expiry(clock, ttl_secs).expect_err("TTL is too short");
3259
3260        assert!(
3261            err.to_string().contains(reason_fragment),
3262            "unexpected error: {err}",
3263        );
3264    }
3265
3266    #[rstest]
3267    #[case(i64::MAX as u64, "overflows Derive signature_expiry_sec")]
3268    #[case(u64::MAX, "cannot fit in Derive signature_expiry_sec")]
3269    fn test_normal_order_signature_expiry_rejects_extreme_ttl(
3270        #[case] ttl_secs: u64,
3271        #[case] reason_fragment: &str,
3272    ) {
3273        let clock = get_atomic_clock_realtime();
3274
3275        let err = normal_order_signature_expiry(clock, ttl_secs).expect_err("TTL is invalid");
3276
3277        assert!(
3278            err.to_string().contains(reason_fragment),
3279            "unexpected error: {err}",
3280        );
3281    }
3282
3283    #[rstest]
3284    #[case(None, "max_fee_per_contract is required")]
3285    #[case(Some(dec!(0)), "max_fee_per_contract must be greater than zero")]
3286    #[case(Some(dec!(-1)), "max_fee_per_contract must be greater than zero")]
3287    fn test_new_rejects_invalid_max_fee_per_contract(
3288        #[case] max_fee_per_contract: Option<Decimal>,
3289        #[case] expected: &str,
3290    ) {
3291        let mut config = test_config();
3292        config.max_fee_per_contract = max_fee_per_contract;
3293
3294        let err = DeriveExecutionClient::new(test_core(), config).expect_err("must reject");
3295
3296        assert_eq!(err.to_string(), expected);
3297    }
3298
3299    #[rstest]
3300    #[case(OrderType::StopMarket, true)]
3301    #[case(OrderType::StopLimit, true)]
3302    #[case(OrderType::MarketIfTouched, true)]
3303    #[case(OrderType::LimitIfTouched, true)]
3304    #[case(OrderType::Market, false)]
3305    #[case(OrderType::Limit, false)]
3306    #[case(OrderType::MarketToLimit, false)]
3307    #[case(OrderType::TrailingStopMarket, false)]
3308    fn test_is_derive_trigger_order_type(#[case] order_type: OrderType, #[case] expected: bool) {
3309        assert_eq!(is_derive_trigger_order_type(order_type), expected);
3310    }
3311
3312    #[rstest]
3313    fn test_resolve_submit_nonce_emits_rejection_and_forgets_identity() {
3314        let clock = get_atomic_clock_realtime();
3315        let instrument_id = InstrumentId::from("ETH-PERP.DERIVE");
3316        let strategy_id = StrategyId::from("S-1");
3317        let client_order_id = ClientOrderId::from("NONCE-SUBMIT-1");
3318        let order = OrderTestBuilder::new(OrderType::Limit)
3319            .trader_id(TraderId::from("TRADER-001"))
3320            .strategy_id(strategy_id)
3321            .instrument_id(instrument_id)
3322            .client_order_id(client_order_id)
3323            .side(OrderSide::Buy)
3324            .quantity(Quantity::from("1.000"))
3325            .price(Price::from("3500.00"))
3326            .build();
3327        let identity = OrderIdentity {
3328            instrument_id,
3329            strategy_id,
3330            order_side: OrderSide::Buy,
3331            order_type: OrderType::Limit,
3332        };
3333        let state = WsDispatchState::new();
3334        state.register_identity(client_order_id, identity);
3335        let (emitter, mut rx) = test_emitter(clock);
3336
3337        let nonce = resolve_submit_nonce(
3338            Err(NonceError::ClockBeforeEpoch),
3339            &emitter,
3340            &state,
3341            &order,
3342            clock,
3343        );
3344        let event = rx.try_recv().expect("OrderRejected event");
3345
3346        assert!(nonce.is_none());
3347        assert!(state.identity(&client_order_id).is_none());
3348        if let ExecutionEvent::Order(OrderEventAny::Rejected(rejected)) = event {
3349            assert_eq!(rejected.client_order_id, client_order_id);
3350            assert_eq!(
3351                rejected.reason.as_str(),
3352                "nonce allocation failed: system clock is before UNIX epoch",
3353            );
3354        } else {
3355            panic!("expected OrderRejected, event was {event:?}");
3356        }
3357    }
3358
3359    #[rstest]
3360    fn test_resolve_modify_nonce_emits_modify_rejection() {
3361        let clock = get_atomic_clock_realtime();
3362        let instrument_id = InstrumentId::from("ETH-PERP.DERIVE");
3363        let strategy_id = StrategyId::from("S-1");
3364        let client_order_id = ClientOrderId::from("NONCE-MODIFY-1");
3365        let venue_order_id = VenueOrderId::from("ord-nonce-modify-1");
3366        let (emitter, mut rx) = test_emitter(clock);
3367
3368        let nonce = resolve_modify_nonce(
3369            Err(NonceError::ClockBeforeEpoch),
3370            &emitter,
3371            strategy_id,
3372            instrument_id,
3373            client_order_id,
3374            venue_order_id,
3375            clock,
3376        );
3377        let event = rx.try_recv().expect("OrderModifyRejected event");
3378
3379        assert!(nonce.is_none());
3380
3381        if let ExecutionEvent::Order(OrderEventAny::ModifyRejected(rejected)) = event {
3382            assert_eq!(rejected.client_order_id, client_order_id);
3383            assert_eq!(rejected.venue_order_id, Some(venue_order_id));
3384            assert_eq!(
3385                rejected.reason.as_str(),
3386                "nonce allocation failed: system clock is before UNIX epoch",
3387            );
3388        } else {
3389            panic!("expected OrderModifyRejected, event was {event:?}");
3390        }
3391    }
3392
3393    #[rstest]
3394    #[case(dec!(0))]
3395    #[case(dec!(-1))]
3396    fn test_round_to_tick_treats_non_positive_tick_as_no_op(#[case] tick: Decimal) {
3397        // Non-positive tick must pass through both sides untouched so the
3398        // signing path does not divide by zero or amplify garbage tick data.
3399        assert_eq!(
3400            round_to_tick(dec!(3501.55), tick, OrderSide::Buy),
3401            dec!(3501.55)
3402        );
3403        assert_eq!(
3404            round_to_tick(dec!(3501.55), tick, OrderSide::Sell),
3405            dec!(3501.55)
3406        );
3407    }
3408
3409    #[rstest]
3410    fn test_resolve_signing_context_rejects_placeholder_domain_separator() {
3411        // The shipped mainnet defaults are real Protocol Constants, so force
3412        // an explicit placeholder via the config override to verify the
3413        // placeholder-detection path still refuses to construct.
3414        let mut config = test_config();
3415        config.environment = DeriveEnvironment::Mainnet;
3416        config.domain_separator =
3417            Some("0x<paste_from_docs.derive.xyz_protocol_constants>".to_string());
3418        let err = DeriveExecutionClient::new(test_core(), config).expect_err("must reject");
3419        let msg = err.to_string();
3420        assert!(msg.contains("placeholder"), "unexpected error: {msg}",);
3421    }
3422
3423    #[rstest]
3424    fn test_resolve_signing_context_uses_mainnet_defaults() {
3425        let mut config = test_config();
3426        config.environment = DeriveEnvironment::Mainnet;
3427        config.domain_separator = None;
3428        config.action_typehash = None;
3429        config.trade_module_address = None;
3430
3431        DeriveExecutionClient::new(test_core(), config).expect("mainnet defaults should parse");
3432    }
3433
3434    #[rstest]
3435    fn test_resolve_signing_context_uses_testnet_defaults() {
3436        let mut config = test_config();
3437        config.environment = DeriveEnvironment::Testnet;
3438        config.domain_separator = None;
3439        config.action_typehash = None;
3440        config.trade_module_address = None;
3441
3442        DeriveExecutionClient::new(test_core(), config).expect("testnet defaults should parse");
3443    }
3444
3445    #[rstest]
3446    fn test_market_order_limit_price_rounds_to_coarse_tick() {
3447        // Coarse tick = 1.0 (e.g. weekly option strikes); raw 3518.505 rounds
3448        // up to 3519, raw 3482.5 rounds down to 3482.
3449        let quote = QuoteTick::new(
3450            InstrumentId::from("ETH-20260627-3500-C.DERIVE"),
3451            Price::from("3500"),
3452            Price::from("3501"),
3453            Quantity::from("1.000"),
3454            Quantity::from("1.000"),
3455            UnixNanos::from(0),
3456            UnixNanos::from(0),
3457        );
3458        let buy = market_order_limit_price(&quote, OrderSide::Buy, 50, dec!(1)).unwrap();
3459        assert_eq!(buy, dec!(3519));
3460        let sell = market_order_limit_price(&quote, OrderSide::Sell, 50, dec!(1)).unwrap();
3461        assert_eq!(sell, dec!(3482));
3462    }
3463
3464    #[rstest]
3465    fn test_new_populates_identity() {
3466        let core = test_core();
3467        let client = DeriveExecutionClient::new(core, test_config()).unwrap();
3468
3469        assert_eq!(client.client_id(), ClientId::from(DERIVE));
3470        assert_eq!(client.account_id(), AccountId::from("DERIVE-001"));
3471        assert_eq!(client.venue(), *DERIVE_VENUE);
3472        assert_eq!(client.oms_type(), OmsType::Netting);
3473        assert_eq!(client.subaccount_id(), TEST_SUBACCOUNT);
3474        assert!(!client.is_connected());
3475    }
3476
3477    #[rstest]
3478    fn test_cache_instrument_registers_report_precision() {
3479        let client = DeriveExecutionClient::new(test_core(), test_config()).unwrap();
3480        let instrument = sample_derive_instrument();
3481        let instrument_id = format_instrument_id(instrument.instrument_name.as_str());
3482
3483        client.cache_instrument(instrument);
3484
3485        assert_eq!(
3486            client.dispatch_state.instrument_precision(&instrument_id),
3487            Some((2, 3)),
3488        );
3489    }
3490
3491    #[rstest]
3492    fn test_order_dispatch_uses_registered_instrument_precision() {
3493        let clock = get_atomic_clock_realtime();
3494        let client = test_client_with_instrument();
3495
3496        let mut order: DeriveOrder = serde_json::from_str(include_str!(
3497            "../test_data/perps/http_order_eth_partially_filled.json"
3498        ))
3499        .unwrap();
3500        order.amount = Decimal::from_str_exact("25.000").unwrap();
3501        order.filled_amount = Decimal::from_str_exact("5.000").unwrap();
3502        order.limit_price = Decimal::from_str_exact("25.000").unwrap();
3503        order.order_status = DeriveOrderStatus::Open;
3504        order.order_type = DeriveOrderType::Limit;
3505
3506        let (emitter, mut rx) = test_emitter(clock);
3507        dispatch_orders_payload(
3508            DeriveOrdersSubscriptionData {
3509                orders: vec![order],
3510            },
3511            &emitter,
3512            AccountId::from("DERIVE-001"),
3513            clock,
3514            &client.dispatch_state,
3515        );
3516
3517        let event = rx.try_recv().unwrap();
3518        let ExecutionEvent::Report(ExecutionReport::Order(report)) = event else {
3519            panic!("Expected OrderStatusReport");
3520        };
3521        assert_eq!(report.price, Some(Price::from("25.00")));
3522        assert_eq!(report.price.unwrap().precision, 2);
3523        assert_eq!(report.quantity, Quantity::from("25.000"));
3524        assert_eq!(report.quantity.precision, 3);
3525        assert_eq!(report.filled_qty, Quantity::from("5.000"));
3526        assert_eq!(report.filled_qty.precision, 3);
3527    }
3528
3529    #[rstest]
3530    fn test_trade_dispatch_uses_registered_instrument_precision() {
3531        let clock = get_atomic_clock_realtime();
3532        let client = test_client_with_instrument();
3533        let mut trade: DeriveTrade = serde_json::from_str(include_str!(
3534            "../test_data/perps/http_private_trade_eth.json"
3535        ))
3536        .unwrap();
3537        trade.trade_amount = Decimal::from_str_exact("25.000").unwrap();
3538        trade.trade_price = Decimal::from_str_exact("25.000").unwrap();
3539
3540        let (emitter, mut rx) = test_emitter(clock);
3541        dispatch_trades_payload(
3542            DeriveTradesSubscriptionData {
3543                trades: vec![trade],
3544            },
3545            &emitter,
3546            AccountId::from("DERIVE-001"),
3547            clock,
3548            &client.dispatch_state,
3549        );
3550
3551        let event = rx.try_recv().unwrap();
3552        let ExecutionEvent::Report(ExecutionReport::Fill(report)) = event else {
3553            panic!("Expected FillReport");
3554        };
3555        assert_eq!(report.last_px, Price::from("25.00"));
3556        assert_eq!(report.last_px.precision, 2);
3557        assert_eq!(report.last_qty, Quantity::from("25.000"));
3558        assert_eq!(report.last_qty.precision, 3);
3559    }
3560
3561    #[rstest]
3562    fn test_emit_tracked_event_suppresses_in_flight_replace_cancel_leg() {
3563        // Derive's `private/replace` cancels the old order; the `.orders`
3564        // cancel-of-old leg can arrive before `modify_order` rebinds the order,
3565        // i.e. while the replace is in flight. In that window only the
3566        // `pending_modify` marker (not the bound-id check) can suppress it. The
3567        // integration suite covers the post-rebind bound-id branch; this covers
3568        // the in-flight branch, which is otherwise unexercised end to end.
3569        let clock = get_atomic_clock_realtime();
3570        let account_id = AccountId::from("DERIVE-001");
3571        let instrument_id = InstrumentId::from("ETH-PERP.DERIVE");
3572        let cid = ClientOrderId::from("STRAT-MOD-INFLIGHT");
3573        let stale_voi = VenueOrderId::from("ord-stale-1");
3574        let identity = OrderIdentity {
3575            instrument_id,
3576            strategy_id: StrategyId::from("S-1"),
3577            order_side: OrderSide::Buy,
3578            order_type: OrderType::Limit,
3579        };
3580        // A `cancelled` report for the stale leg, identical across both cases:
3581        // only the dispatch-state marker differs.
3582        let report = OrderStatusReport::new(
3583            account_id,
3584            instrument_id,
3585            Some(cid),
3586            stale_voi,
3587            OrderSide::Buy.into(),
3588            OrderType::Limit,
3589            TimeInForce::Gtc,
3590            OrderStatus::Canceled,
3591            Quantity::from("1.000"),
3592            Quantity::from("0.000"),
3593            UnixNanos::from(1_000),
3594            UnixNanos::from(2_000),
3595            UnixNanos::from(3_000),
3596            None,
3597        );
3598
3599        // Marker targets the cancel's venue order id and no bound id is
3600        // recorded, so suppression can only come from the in-flight branch.
3601        let (emitter, mut rx) = test_emitter(clock);
3602        let state = WsDispatchState::new();
3603        state.mark_pending_modify(cid, stale_voi);
3604        emit_tracked_order_event(
3605            &emitter,
3606            &state,
3607            cid,
3608            identity,
3609            &report,
3610            account_id,
3611            UnixNanos::from(0),
3612        );
3613        let suppressed = rx.try_recv().is_err();
3614
3615        // A marker for a different venue order id must not suppress: the guard
3616        // keys on the specific id, so the cancel-of-old still terminates.
3617        let (emitter, mut rx) = test_emitter(clock);
3618        let state = WsDispatchState::new();
3619        state.mark_pending_modify(cid, VenueOrderId::from("ord-other"));
3620        emit_tracked_order_event(
3621            &emitter,
3622            &state,
3623            cid,
3624            identity,
3625            &report,
3626            account_id,
3627            UnixNanos::from(0),
3628        );
3629        let mut saw_canceled = false;
3630
3631        while let Ok(event) = rx.try_recv() {
3632            if matches!(event, ExecutionEvent::Order(OrderEventAny::Canceled(_))) {
3633                saw_canceled = true;
3634            }
3635        }
3636
3637        assert!(
3638            suppressed,
3639            "in-flight cancel-of-old leg must be suppressed by the pending-modify marker",
3640        );
3641        assert!(
3642            saw_canceled,
3643            "a pending-modify marker for a different venue order id must not suppress",
3644        );
3645    }
3646
3647    #[rstest]
3648    fn test_ensure_canceled_emitted_is_idempotent() {
3649        let clock = get_atomic_clock_realtime();
3650        let account_id = AccountId::from("DERIVE-001");
3651        let client_order_id = ClientOrderId::from("TRIGGER-CANCEL-1");
3652        let identity = OrderIdentity {
3653            instrument_id: InstrumentId::from("ETH-PERP.DERIVE"),
3654            strategy_id: StrategyId::from("S-1"),
3655            order_side: OrderSide::Buy,
3656            order_type: OrderType::StopMarket,
3657        };
3658        let venue_order_id = VenueOrderId::from("trigger-cancel-1");
3659        let state = WsDispatchState::new();
3660        let (emitter, mut rx) = test_emitter(clock);
3661
3662        for _ in 0..2 {
3663            ensure_canceled_emitted(
3664                &emitter,
3665                &state,
3666                client_order_id,
3667                identity,
3668                venue_order_id,
3669                account_id,
3670                UnixNanos::from(1_000),
3671                UnixNanos::from(1_000),
3672            );
3673        }
3674
3675        assert!(matches!(
3676            rx.try_recv(),
3677            Ok(ExecutionEvent::Order(OrderEventAny::Canceled(_)))
3678        ));
3679        assert!(rx.try_recv().is_err(), "duplicate OrderCanceled emitted");
3680    }
3681
3682    fn test_emitter(
3683        clock: &'static AtomicTime,
3684    ) -> (
3685        ExecutionEventEmitter,
3686        tokio::sync::mpsc::UnboundedReceiver<ExecutionEvent>,
3687    ) {
3688        let (tx, rx) = tokio::sync::mpsc::unbounded_channel();
3689        let mut emitter = ExecutionEventEmitter::new(
3690            clock,
3691            TraderId::from("TRADER-001"),
3692            AccountId::from("DERIVE-001"),
3693            AccountType::Margin,
3694            Some(Currency::USDC()),
3695        );
3696        emitter.set_sender(tx);
3697        (emitter, rx)
3698    }
3699
3700    fn sample_derive_instrument() -> DeriveInstrument {
3701        serde_json::from_str(include_str!("../test_data/perps/instrument_eth.json")).unwrap()
3702    }
3703
3704    fn test_client_with_instrument() -> DeriveExecutionClient {
3705        let mut client = DeriveExecutionClient::new(test_core(), test_config()).unwrap();
3706        let instrument =
3707            parse_derive_instrument_any(&sample_derive_instrument(), UnixNanos::default())
3708                .unwrap()
3709                .unwrap();
3710        client.on_instrument(instrument);
3711        client
3712    }
3713}