1use 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#[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 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 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 #[must_use]
246 pub const fn subaccount_id(&self) -> u64 {
247 self.credential.subaccount_id()
248 }
249
250 #[must_use]
252 pub fn config(&self) -> &DeriveExecClientConfig {
253 &self.config
254 }
255
256 #[must_use]
258 pub fn http_client(&self) -> &DeriveHttpClient {
259 &self.http_client
260 }
261
262 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 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 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 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 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 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 }
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 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(¶ms).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 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 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 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 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 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 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 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 "e,
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 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 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 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 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 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 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(¶ms).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(¶ms).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 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
2329fn 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
2421pub 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
2462pub 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#[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 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 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
2793fn 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 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
2957fn 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 let price = market_order_limit_price("e, 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 let price = market_order_limit_price("e, OrderSide::Sell, 50, dec!(0.01)).unwrap();
3077 assert_eq!(price, dec!(3482.5));
3078
3079 let zero = market_order_limit_price("e, 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 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 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 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("e, OrderSide::Buy, 50, dec!(1)).unwrap();
3316 assert_eq!(buy, dec!(3519));
3317 let sell = market_order_limit_price("e, 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 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 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 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 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}