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