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