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