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