1use std::{
19 future::Future,
20 time::{Duration, Instant},
21};
22
23use anyhow::Context;
24use async_trait::async_trait;
25use futures_util::{StreamExt, pin_mut};
26use nautilus_common::{
27 clients::ExecutionClient,
28 live::runner::get_exec_event_sender,
29 messages::execution::{
30 BatchCancelOrders, CancelAllOrders, CancelOrder, GenerateFillReports,
31 GenerateOrderStatusReport, GenerateOrderStatusReports, GeneratePositionStatusReports,
32 ModifyOrder, QueryAccount, QueryOrder, SubmitOrder, SubmitOrderList,
33 },
34};
35use nautilus_core::{
36 AtomicMap, Params, UUID4, UnixNanos,
37 time::{AtomicTime, get_atomic_clock_realtime},
38};
39use nautilus_live::{
40 ExecutionClientCore, ExecutionEventEmitter, SocketControl,
41 execution::failure::CommandFailure,
42 task::{TaskGroup, TaskGroupGuard},
43};
44use nautilus_model::{
45 accounts::AccountAny,
46 enums::{AccountType, LiquiditySide, OmsType, OrderSide, OrderStatus, OrderType, TimeInForce},
47 events::{
48 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderEventAny, OrderExpired,
49 OrderFilled, OrderInitialized, OrderRejected, OrderUpdated,
50 },
51 identifiers::{
52 AccountId, ClientId, ClientOrderId, InstrumentId, StrategyId, TradeId, Venue, VenueOrderId,
53 },
54 instruments::{Instrument, InstrumentAny},
55 orders::{Order, OrderAny},
56 reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
57 types::{AccountBalance, MarginBalance, Money, Price, Quantity},
58};
59use ustr::Ustr;
60
61use crate::{
62 common::{
63 auth::run_auth_token_refresh,
64 consts::{
65 AX_ACCOUNT_REGISTRATION_TIMEOUT_SECS, AX_AUTH_TOKEN_TTL_SECS, AX_POST_ONLY_REJECT,
66 AX_VENUE,
67 },
68 credential::Credential,
69 enums::{AxOrderSide, AxTimeInForce},
70 parse::{ax_timestamp_stn_to_unix_nanos, cid_to_client_order_id, quantity_to_contracts},
71 },
72 config::AxExecutionClientConfig,
73 http::{
74 client::AxHttpClient,
75 error::AxHttpError,
76 models::{AxOrderRejectReason, PreviewAggressiveLimitOrderRequest, ReplaceOrderRequest},
77 },
78 websocket::{
79 AxOrdersWsMessage, AxWsOrderEvent,
80 messages::{AxWsOrder, AxWsTradeExecution, OrderMetadata},
81 orders::{AxOrdersWebSocketClient, AxOrdersWsClientError, OrdersCaches},
82 },
83};
84
85#[derive(Debug)]
87pub struct AxExecutionClient {
88 core: ExecutionClientCore,
89 clock: &'static AtomicTime,
90 config: AxExecutionClientConfig,
91 emitter: ExecutionEventEmitter,
92 http_client: AxHttpClient,
93 ws_orders: AxOrdersWebSocketClient,
94 session_tasks: TaskGroup,
95 pending_tasks: TaskGroup,
96 shutdown_errors: Vec<String>,
97}
98
99impl AxExecutionClient {
100 pub fn new(core: ExecutionClientCore, config: AxExecutionClientConfig) -> anyhow::Result<Self> {
106 let http_client = AxHttpClient::with_credentials(
107 config.api_key.clone().unwrap_or_default(),
108 config.api_secret.clone().unwrap_or_default(),
109 Some(config.http_base_url()),
110 Some(config.orders_base_url()),
111 config.http_timeout_secs,
112 config.max_retries,
113 config.retry_delay_initial_ms,
114 config.retry_delay_max_ms,
115 config.proxy_url.clone(),
116 )?;
117
118 let clock = get_atomic_clock_realtime();
119 let trader_id = core.trader_id;
120 let account_id = core.account_id;
121 let emitter =
122 ExecutionEventEmitter::new(clock, trader_id, account_id, AccountType::Margin, None);
123 let mut ws_url = config.ws_private_url();
124 if config.cancel_on_disconnect {
125 let separator = if ws_url.contains('?') { "&" } else { "?" };
126 ws_url.push_str(&format!("{separator}cancel_on_disconnect=true"));
127 }
128 let ws_orders = AxOrdersWebSocketClient::new(
129 ws_url,
130 account_id,
131 trader_id,
132 config.heartbeat_interval_secs,
133 config.transport_backend,
134 config.proxy_url.clone(),
135 )
136 .with_socket_control(SocketControl::new(
137 core.client_id,
138 Some(*AX_VENUE),
139 "architect-ax-user-streams",
140 ));
141
142 let session_tasks = TaskGroup::new();
143 let pending_tasks = TaskGroup::new();
144
145 Ok(Self {
146 core,
147 clock,
148 config,
149 emitter,
150 http_client,
151 ws_orders,
152 session_tasks,
153 pending_tasks,
154 shutdown_errors: Vec::new(),
155 })
156 }
157
158 async fn authenticate(&self, credential: &Credential) -> anyhow::Result<String> {
159 self.http_client
160 .authenticate(
161 credential.api_key(),
162 credential.api_secret(),
163 AX_AUTH_TOKEN_TTL_SECS,
164 )
165 .await
166 .map_err(|e| anyhow::anyhow!("Authentication failed: {e}"))
167 }
168
169 fn update_account_state(&self) {
170 let http_client = self.http_client.clone();
171 let account_id = self.core.account_id;
172 let emitter = self.emitter.clone();
173 let clock = self.clock;
174
175 self.spawn_task("query_account", async move {
176 let account_state = http_client
177 .request_account_state(account_id)
178 .await
179 .context("failed to request AX account state")?;
180 let ts_event = clock.get_time_ns();
181 emitter.emit_account_state(
182 account_state.balances.clone(),
183 account_state.margins.clone(),
184 account_state.is_reported,
185 ts_event,
186 account_state.info,
187 );
188 Ok(())
189 });
190 }
191
192 fn submit_order_internal(&self, cmd: &SubmitOrder) -> anyhow::Result<()> {
193 let (
194 order_for_task,
195 client_order_id,
196 strategy_id,
197 instrument_id,
198 order_side,
199 order_type,
200 quantity,
201 time_in_force,
202 is_post_only,
203 limit_price,
204 ) = {
205 let cache = self.core.cache();
206 let order = cache.try_order(&cmd.client_order_id)?;
207 (
208 order.clone(),
209 order.client_order_id(),
210 order.strategy_id(),
211 order.instrument_id(),
212 order.order_side(),
213 order.order_type(),
214 order.quantity(),
215 order.time_in_force(),
216 order.is_post_only(),
217 order.price(),
218 )
219 };
220
221 let ws_orders = self.ws_orders.clone();
222 let trader_id = self.core.trader_id;
223 let emitter = self.emitter.clone();
224 let clock = self.clock;
225
226 let http_client = self.http_client.clone();
227
228 self.spawn_task("submit_order", async move {
229 let (price, submit_time_in_force, submit_post_only) = if order_type
232 == OrderType::Market
233 {
234 let preview_result: anyhow::Result<Price> = async {
235 let symbol = instrument_id.symbol.inner();
236 let ax_side = AxOrderSide::from(order_side);
237 let qty_contracts = quantity_to_contracts(quantity)?;
238
239 let instrument = http_client.get_instrument(&symbol).ok_or_else(|| {
240 anyhow::anyhow!("Instrument {instrument_id} not found in cache")
241 })?;
242
243 let request =
244 PreviewAggressiveLimitOrderRequest::new(symbol, qty_contracts, ax_side);
245 let response = http_client
246 .inner
247 .preview_aggressive_limit_order(&request)
248 .await
249 .map_err(|e| {
250 anyhow::anyhow!("Failed to preview aggressive limit order: {e}")
251 })?;
252
253 if response.remaining_quantity > 0 {
254 log::warn!(
255 "Market order book depth insufficient: \
256 filled_qty={} remaining_qty={} for {instrument_id}",
257 response.filled_quantity,
258 response.remaining_quantity,
259 );
260 }
261
262 let limit_price_decimal = response.limit_price.ok_or_else(|| {
263 anyhow::anyhow!(
264 "No liquidity available for market order on {instrument_id}"
265 )
266 })?;
267
268 let price =
269 Price::from_decimal_dp(limit_price_decimal, instrument.price_precision())
270 .with_context(|| {
271 format!(
272 "Failed to convert AX take-through price {limit_price_decimal} for {instrument_id}"
273 )
274 })?;
275 log::debug!("Market order take-through price: {price} for {instrument_id}",);
276 Ok(price)
277 }
278 .await;
279
280 let price = match preview_result {
281 Ok(price) => price,
282 Err(e) => {
283 let reason = e.to_string();
284 log::warn!(
285 "AX market order preview failed for {client_order_id}: {reason}"
286 );
287 emitter.emit_order_rejected(
288 &order_for_task,
289 &reason,
290 clock.get_time_ns(),
291 false,
292 );
293 return Ok(());
294 }
295 };
296
297 (price, TimeInForce::Ioc, false)
298 } else {
299 (
300 limit_price.context("AX limit order is missing a price")?,
301 time_in_force,
302 is_post_only,
303 )
304 };
305
306 let result = ws_orders
307 .submit_order(
308 trader_id,
309 strategy_id,
310 instrument_id,
311 client_order_id,
312 order_side,
313 quantity,
314 submit_time_in_force,
315 price,
316 submit_post_only,
317 )
318 .await;
319
320 if let Err(e) = result {
321 match classify_ax_ws_failure(&e) {
322 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
325 log::warn!("AX submit failed for {client_order_id}: {reason}");
326 emitter.emit_order_rejected(
327 &order_for_task,
328 &reason,
329 clock.get_time_ns(),
330 false,
331 );
332 }
333 CommandFailure::Ambiguous(reason) => {
334 log::warn!(
335 "Ambiguous AX submit failure for {client_order_id}, awaiting reconciliation: {reason}"
336 );
337 }
338 }
339 }
340
341 Ok(())
342 });
343
344 Ok(())
345 }
346
347 fn cancel_order_internal(&self, cmd: &CancelOrder) {
348 let ws_orders = self.ws_orders.clone();
349 let client_order_id = cmd.client_order_id;
350 let venue_order_id = cmd.venue_order_id;
351
352 self.spawn_task("cancel_order", async move {
355 if let Err(e) = ws_orders.cancel_order(client_order_id, venue_order_id).await {
356 match classify_ax_ws_failure(&e) {
357 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
358 log::warn!("Cancel command failed for {client_order_id}: {reason}");
359 }
360 CommandFailure::Ambiguous(reason) => {
361 log::warn!(
362 "Ambiguous AX cancel failure for {client_order_id}, awaiting reconciliation: {reason}"
363 );
364 }
365 }
366 }
367
368 Ok(())
369 });
370 }
371
372 fn spawn_task<F>(&self, description: &'static str, fut: F)
373 where
374 F: Future<Output = anyhow::Result<()>> + Send + 'static,
375 {
376 let future = async move {
377 if let Err(e) = fut.await {
378 log::warn!("{description} failed: {e}");
379 }
380 };
381
382 if let Err(e) = self.pending_tasks.spawn(future) {
383 log::warn!("Skipping AX {description} after shutdown began: {e}");
384 }
385 }
386
387 fn abort_pending_tasks(&self) {
388 self.pending_tasks.begin_shutdown();
389 }
390
391 fn abort_session_tasks(&self) {
392 self.session_tasks.begin_shutdown();
393 self.ws_orders.begin_shutdown();
394 }
395
396 async fn await_pending_tasks(&self) -> anyhow::Result<()> {
397 self.pending_tasks.begin_shutdown();
398 self.pending_tasks
399 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
400 .await
401 .map_err(|e| anyhow::anyhow!("Failed to terminate AX execution tasks: {e}"))?;
402 Ok(())
403 }
404
405 async fn await_session_tasks(&self) -> anyhow::Result<()> {
406 self.session_tasks.begin_shutdown();
407 self.session_tasks
408 .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2))
409 .await
410 .map_err(|e| anyhow::anyhow!("Failed to terminate AX execution session tasks: {e}"))?;
411 Ok(())
412 }
413
414 async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
415 self.abort_session_tasks();
416 self.abort_pending_tasks();
417 self.http_client.cancel_all_requests();
418
419 if let Err(e) = self.ws_orders.close().await {
420 self.shutdown_errors
421 .push(format!("AX orders WebSocket shutdown failed: {e}"));
422 }
423
424 let (session_result, pending_result) =
425 tokio::join!(self.await_session_tasks(), self.await_pending_tasks());
426 self.core.set_disconnected();
427
428 if let Err(e) = session_result {
429 self.shutdown_errors.push(e.to_string());
430 }
431
432 if let Err(e) = pending_result {
433 self.shutdown_errors.push(e.to_string());
434 }
435
436 if !self.shutdown_errors.is_empty() {
437 anyhow::bail!(std::mem::take(&mut self.shutdown_errors).join("; "));
438 }
439 Ok(())
440 }
441
442 async fn await_account_registered(&self, timeout_secs: f64) -> anyhow::Result<()> {
444 let account_id = self.core.account_id;
445
446 if self.core.cache().account(&account_id).is_some() {
447 log::info!("Account {account_id} registered");
448 return Ok(());
449 }
450
451 let start = Instant::now();
452 let timeout = Duration::from_secs_f64(timeout_secs);
453 let interval = Duration::from_millis(10);
454
455 loop {
456 tokio::time::sleep(interval).await;
457
458 if self.core.cache().account(&account_id).is_some() {
459 log::info!("Account {account_id} registered");
460 return Ok(());
461 }
462
463 if start.elapsed() >= timeout {
464 anyhow::bail!(
465 "Timeout waiting for account {account_id} to be registered after {timeout_secs}s"
466 );
467 }
468 }
469 }
470}
471
472#[async_trait(?Send)]
473impl ExecutionClient for AxExecutionClient {
474 fn is_connected(&self) -> bool {
475 self.core.is_connected()
476 }
477
478 fn client_id(&self) -> ClientId {
479 self.core.client_id
480 }
481
482 fn account_id(&self) -> AccountId {
483 self.core.account_id
484 }
485
486 fn venue(&self) -> Venue {
487 *AX_VENUE
488 }
489
490 fn oms_type(&self) -> OmsType {
491 self.core.oms_type
492 }
493
494 fn get_account(&self) -> Option<AccountAny> {
495 self.core.cache().account_owned(&self.core.account_id)
496 }
497
498 async fn connect(&mut self) -> anyhow::Result<()> {
499 if self.core.is_connected() && self.pending_tasks.is_open() && self.session_tasks.is_open()
500 {
501 return Ok(());
502 }
503
504 if !self.pending_tasks.is_open() || !self.session_tasks.is_open() {
505 self.teardown_partial_connect().await?;
506 self.pending_tasks.start_generation().map_err(|e| {
507 anyhow::anyhow!("Failed to start AX execution task generation: {e}")
508 })?;
509 self.session_tasks.start_generation().map_err(|e| {
510 anyhow::anyhow!("Failed to start AX execution session generation: {e}")
511 })?;
512 }
513 let http_client = self.http_client.clone();
514 let ws_orders = self.ws_orders.clone();
515 let setup_guard =
516 TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
517 http_client.cancel_all_requests();
518 ws_orders.begin_shutdown();
519 });
520
521 self.http_client.reset_cancellation_token();
523
524 let credential =
525 Credential::resolve(self.config.api_key.clone(), self.config.api_secret.clone())
526 .context("API credentials not configured")?;
527 let token = self.authenticate(&credential).await?;
528
529 if !self.core.instruments_initialized() {
533 self.http_client
534 .request_account_fees()
535 .await
536 .context("failed to resolve AX account fee rates")?;
537
538 let instruments = self
539 .http_client
540 .request_instruments(None, None)
541 .await
542 .context("failed to request AX instruments")?;
543
544 if instruments.is_empty() {
545 log::warn!("No instruments returned from AX");
546 } else {
547 log::debug!("Loaded {} instruments", instruments.len());
548 self.http_client.cache_instruments(&instruments);
549 self.ws_orders.cache_instruments(&instruments);
550 }
551 self.core.set_instruments_initialized();
552 }
553
554 self.ws_orders.connect(&token).await?;
555 log::debug!("Connected to orders WebSocket");
556
557 let stream = self.ws_orders.stream();
558 let emitter = self.emitter.clone();
559 let caches = self.ws_orders.caches().clone();
560 let account_id = self.core.account_id;
561 let instruments_cache = self.ws_orders.instruments_cache();
562 let clock = self.clock;
563
564 if let Err(e) = self.session_tasks.spawn(async move {
565 pin_mut!(stream);
566 while let Some(message) = stream.next().await {
567 dispatch_ws_message(
568 message,
569 &emitter,
570 &caches,
571 account_id,
572 &instruments_cache,
573 clock,
574 );
575 }
576 }) {
577 if let Err(teardown_error) = self.teardown_partial_connect().await {
578 return Err(anyhow::Error::new(e).context(format!(
579 "AX execution startup teardown failed: {teardown_error}"
580 )));
581 }
582 return Err(e.into());
583 }
584
585 let session_result = async {
586 let account_state = self
587 .http_client
588 .request_account_state(self.core.account_id)
589 .await
590 .context("failed to request AX account state")?;
591
592 if !account_state.balances.is_empty() {
593 log::debug!(
594 "Received account state with {} balance(s)",
595 account_state.balances.len()
596 );
597 }
598 self.emitter.send_account_state(account_state);
599
600 self.await_account_registered(AX_ACCOUNT_REGISTRATION_TIMEOUT_SECS)
601 .await?;
602
603 let ws_orders = self.ws_orders.clone();
604 self.session_tasks.spawn(run_auth_token_refresh(
605 self.http_client.clone(),
606 credential,
607 move |token| ws_orders.update_auth_token(&token),
608 ))?;
609 Ok::<(), anyhow::Error>(())
610 }
611 .await;
612
613 if let Err(e) = session_result {
614 if let Err(teardown_error) = self.teardown_partial_connect().await {
615 return Err(e.context(format!(
616 "AX execution startup teardown failed: {teardown_error}"
617 )));
618 }
619 return Err(e);
620 }
621
622 self.core.set_connected();
623 setup_guard.disarm();
624 log::info!("Connected: client_id={}", self.core.client_id);
625 Ok(())
626 }
627
628 async fn disconnect(&mut self) -> anyhow::Result<()> {
629 self.abort_session_tasks();
630 self.abort_pending_tasks();
631 self.http_client.cancel_all_requests();
632
633 let ws_result = self.ws_orders.close().await;
634 let (session_result, pending_result) =
635 tokio::join!(self.await_session_tasks(), self.await_pending_tasks());
636
637 self.core.set_disconnected();
638 ws_result?;
639 session_result?;
640 pending_result?;
641 log::info!("Disconnected: client_id={}", self.core.client_id);
642 Ok(())
643 }
644
645 fn query_account(&self, _cmd: QueryAccount) -> anyhow::Result<()> {
646 self.update_account_state();
647 Ok(())
648 }
649
650 fn query_order(&self, cmd: QueryOrder) -> anyhow::Result<()> {
651 let http_client = self.http_client.clone();
652 let account_id = self.core.account_id;
653 let client_order_id = cmd.client_order_id;
654 let venue_order_id = cmd.venue_order_id.or_else(|| {
655 self.ws_orders
656 .orders_metadata()
657 .get(&client_order_id)
658 .and_then(|metadata| metadata.venue_order_id)
659 });
660 let instrument_id = cmd.instrument_id;
661 let emitter = self.emitter.clone();
662
663 let (order_side, order_type, time_in_force) = {
665 let cache = self.core.cache();
666 match cache.order(&client_order_id) {
667 Some(order) => (
668 Some(order.order_side()),
669 order.order_type(),
670 order.time_in_force(),
671 ),
672 None => (None, OrderType::Limit, TimeInForce::Gtc),
673 }
674 };
675
676 self.spawn_task("query_order", async move {
677 match http_client
678 .request_order_status(
679 account_id,
680 instrument_id,
681 Some(client_order_id),
682 venue_order_id,
683 order_side,
684 order_type,
685 time_in_force,
686 )
687 .await
688 {
689 Ok(report) => emitter.send_order_status_report(report),
690 Err(e) => log::error!("AX query order failed: {e}"),
691 }
692 Ok(())
693 });
694
695 Ok(())
696 }
697
698 fn generate_account_state(
699 &self,
700 balances: Vec<AccountBalance>,
701 margins: Vec<MarginBalance>,
702 reported: bool,
703 ts_event: UnixNanos,
704 info: Option<Params>,
705 ) -> anyhow::Result<()> {
706 self.emitter
707 .emit_account_state(balances, margins, reported, ts_event, info);
708 Ok(())
709 }
710
711 fn start(&mut self) -> anyhow::Result<()> {
712 if self.core.is_started() {
713 return Ok(());
714 }
715
716 self.emitter.set_sender(get_exec_event_sender());
717 self.core.set_started();
718 log::info!(
719 "Started: client_id={}, account_id={}, environment={}",
720 self.core.client_id,
721 self.core.account_id,
722 self.config.environment,
723 );
724 Ok(())
725 }
726
727 fn stop(&mut self) -> anyhow::Result<()> {
728 if self.core.is_stopped() {
729 return Ok(());
730 }
731
732 self.core.set_stopped();
733 self.core.set_disconnected();
734
735 self.abort_session_tasks();
736 self.abort_pending_tasks();
737 log::info!("Stopped: client_id={}", self.core.client_id);
738 Ok(())
739 }
740
741 fn reset(&mut self) -> anyhow::Result<()> {
742 self.abort_session_tasks();
743 self.abort_pending_tasks();
744 self.core.set_disconnected();
745 Ok(())
746 }
747
748 fn dispose(&mut self) -> anyhow::Result<()> {
749 self.abort_session_tasks();
750 self.abort_pending_tasks();
751 self.core.set_disconnected();
752 Ok(())
753 }
754
755 fn submit_order(&self, cmd: SubmitOrder) -> anyhow::Result<()> {
756 {
757 let cache = self.core.cache();
758 let order = cache.try_order(&cmd.client_order_id)?;
759
760 if order.is_closed() {
761 log::warn!("Cannot submit closed order {}", order.client_order_id());
762 return Ok(());
763 }
764
765 if let Err(e) = validate_order_for_ax_submit(&order)
766 .and_then(|()| validate_order_init_instructions(&cmd.order_init))
767 {
768 self.emitter.emit_order_denied(&order, &e.to_string());
769 return Ok(());
770 }
771
772 log::debug!("OrderSubmitted client_order_id={}", order.client_order_id());
773 self.emitter.emit_order_submitted(&order);
774 }
775
776 self.submit_order_internal(&cmd)
777 }
778
779 fn submit_order_list(&self, cmd: SubmitOrderList) -> anyhow::Result<()> {
780 for (client_order_id, order_init) in cmd
781 .order_list
782 .client_order_ids
783 .iter()
784 .zip(cmd.order_inits.iter())
785 {
786 let submit_cmd = SubmitOrder::new(
787 cmd.trader_id,
788 cmd.client_id,
789 cmd.strategy_id,
790 cmd.instrument_id,
791 *client_order_id,
792 order_init.clone(),
793 cmd.exec_algorithm_id,
794 cmd.position_id,
795 cmd.params.clone(),
796 UUID4::new(),
797 cmd.ts_init,
798 cmd.correlation_id,
799 );
800 self.submit_order(submit_cmd)?;
801 }
802 Ok(())
803 }
804
805 fn modify_order(&self, cmd: ModifyOrder) -> anyhow::Result<()> {
806 if cmd.trigger_price.is_some() {
807 emit_ax_modify_rejected(
808 &self.emitter,
809 self.clock,
810 cmd.strategy_id,
811 cmd.instrument_id,
812 cmd.client_order_id,
813 cmd.venue_order_id,
814 "AX does not support venue-native trigger prices",
815 );
816 return Ok(());
817 }
818
819 let venue_order_id = match cmd.venue_order_id {
820 Some(ref voi) => *voi,
821 None => {
822 emit_ax_modify_rejected(
823 &self.emitter,
824 self.clock,
825 cmd.strategy_id,
826 cmd.instrument_id,
827 cmd.client_order_id,
828 None,
829 "missing venue_order_id",
830 );
831 return Ok(());
832 }
833 };
834
835 let quantity = match cmd.quantity {
836 Some(quantity) => match quantity_to_contracts(quantity) {
837 Ok(contracts) => Some(contracts),
838 Err(e) => {
839 emit_ax_modify_rejected(
840 &self.emitter,
841 self.clock,
842 cmd.strategy_id,
843 cmd.instrument_id,
844 cmd.client_order_id,
845 Some(venue_order_id),
846 &e.to_string(),
847 );
848 return Ok(());
849 }
850 },
851 None => None,
852 };
853
854 if !self.core.is_connected() {
855 emit_ax_modify_rejected(
856 &self.emitter,
857 self.clock,
858 cmd.strategy_id,
859 cmd.instrument_id,
860 cmd.client_order_id,
861 Some(venue_order_id),
862 "AX execution client is not connected",
863 );
864 return Ok(());
865 }
866
867 let http_client = self.http_client.clone();
868 let emitter = self.emitter.clone();
869 let clock = self.clock;
870 let strategy_id = cmd.strategy_id;
871 let instrument_id = cmd.instrument_id;
872 let caches = self.ws_orders.caches().clone();
873 let client_order_id = cmd.client_order_id;
874 let price = cmd.price;
875
876 self.spawn_task("modify_order", async move {
877 let mut request = ReplaceOrderRequest::new(venue_order_id.as_str());
878
879 if let Some(price) = price {
880 request = request.with_price(price.as_decimal());
881 }
882
883 if let Some(contracts) = quantity {
884 request = request.with_quantity(contracts);
885 }
886
887 match http_client.inner.replace_order(&request).await {
888 Ok(resp) => {
889 let new_venue_order_id = match VenueOrderId::new_checked(&resp.oid) {
890 Ok(venue_order_id) => venue_order_id,
891 Err(e) => {
892 log::warn!(
893 "AX replace returned invalid venue order ID for {client_order_id}, awaiting reconciliation: {e}"
894 );
895 return Ok(());
896 }
897 };
898 record_replacement_venue_id(
899 &caches,
900 client_order_id,
901 new_venue_order_id,
902 false,
903 );
904 log::debug!("Order replaced: old={} new={}", request.oid, resp.oid);
905 }
906 Err(e) => match classify_ax_http_failure(&e) {
909 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
910 emit_ax_modify_rejected(
911 &emitter,
912 clock,
913 strategy_id,
914 instrument_id,
915 client_order_id,
916 Some(venue_order_id),
917 &reason,
918 );
919 }
920 CommandFailure::Ambiguous(reason) => {
921 log::warn!(
922 "Ambiguous AX modify failure for {client_order_id}, awaiting reconciliation: {reason}"
923 );
924 }
925 },
926 }
927
928 Ok(())
929 });
930
931 Ok(())
932 }
933
934 fn cancel_order(&self, cmd: CancelOrder) -> anyhow::Result<()> {
935 self.cancel_order_internal(&cmd);
936 Ok(())
937 }
938
939 fn cancel_all_orders(&self, cmd: CancelAllOrders) -> anyhow::Result<()> {
940 let http_client = self.http_client.clone();
941 let emitter = self.emitter.clone();
942 let clock = self.clock;
943 let instrument_id = cmd.instrument_id;
944 let account_id = self.core.account_id;
945 let trader_id = self.core.trader_id;
946
947 let open_orders: Vec<(ClientOrderId, Option<VenueOrderId>, StrategyId)> = {
949 let cache = self.core.cache();
950 cache
951 .orders_open(None, Some(&instrument_id), None, None, None)
952 .iter()
953 .map(|o| (o.client_order_id(), o.venue_order_id(), o.strategy_id()))
954 .collect()
955 };
956
957 let caches = self.ws_orders.caches().clone();
958
959 self.spawn_task("cancel_all_orders", async move {
960 match http_client.cancel_all_orders(instrument_id).await {
961 Ok(()) => {
962 log::debug!("Canceled all orders for {instrument_id}");
963
964 let ts_event = clock.get_time_ns();
968
969 for (client_order_id, venue_order_id, strategy_id) in &open_orders {
970 let event = OrderCanceled::new(
971 trader_id,
972 *strategy_id,
973 instrument_id,
974 *client_order_id,
975 UUID4::new(),
976 ts_event,
977 clock.get_time_ns(),
978 false,
979 *venue_order_id,
980 Some(account_id),
981 );
982 emitter.send_order_event(OrderEventAny::Canceled(event));
983
984 if let Some(voi) = venue_order_id {
985 caches.venue_to_client_id.remove(voi);
986 }
987 caches.orders_metadata.remove(client_order_id);
988 caches
989 .cid_to_client_order_id
990 .retain(|_, mapped_client_order_id| {
991 mapped_client_order_id != client_order_id
992 });
993 }
994 }
995 Err(e) => match classify_ax_http_failure(&e) {
998 CommandFailure::NotSent(reason) | CommandFailure::VenueRejected(reason) => {
999 log::warn!("Cancel-all for {instrument_id} failed: {reason}");
1000 }
1001 CommandFailure::Ambiguous(reason) => {
1002 log::warn!(
1003 "Ambiguous AX cancel-all failure for {instrument_id}, awaiting reconciliation: {reason}"
1004 );
1005 }
1006 },
1007 }
1008 Ok(())
1009 });
1010
1011 Ok(())
1012 }
1013
1014 fn batch_cancel_orders(&self, cmd: BatchCancelOrders) -> anyhow::Result<()> {
1015 for cancel in &cmd.cancels {
1016 self.cancel_order_internal(cancel);
1017 }
1018 Ok(())
1019 }
1020
1021 async fn generate_order_status_report(
1022 &self,
1023 cmd: &GenerateOrderStatusReport,
1024 ) -> anyhow::Result<Option<OrderStatusReport>> {
1025 let cid_map = self.ws_orders.cid_to_client_order_id().clone();
1026 let cid_resolver = move |cid: u64| cid_map.get(&cid).map(|v| *v);
1027
1028 let mut reports = self
1029 .http_client
1030 .request_order_status_reports(self.core.account_id, Some(cid_resolver))
1031 .await?;
1032
1033 if let Some(instrument_id) = cmd.instrument_id {
1034 reports.retain(|report| report.instrument_id == instrument_id);
1035 }
1036
1037 if let Some(client_order_id) = cmd.client_order_id {
1038 reports.retain(|report| report.client_order_id == Some(client_order_id));
1039 }
1040
1041 if let Some(venue_order_id) = cmd.venue_order_id {
1042 reports.retain(|report| report.venue_order_id.as_str() == venue_order_id.as_str());
1043 }
1044
1045 Ok(reports.into_iter().next())
1046 }
1047
1048 async fn generate_order_status_reports(
1049 &self,
1050 cmd: &GenerateOrderStatusReports,
1051 ) -> anyhow::Result<Vec<OrderStatusReport>> {
1052 let cid_map = self.ws_orders.cid_to_client_order_id().clone();
1053 let cid_resolver = move |cid: u64| cid_map.get(&cid).map(|v| *v);
1054
1055 let mut reports = if cmd.open_only {
1056 self.http_client
1057 .request_order_status_reports(self.core.account_id, Some(cid_resolver))
1058 .await?
1059 } else {
1060 self.http_client
1061 .request_historical_order_status_reports(
1062 self.core.account_id,
1063 cmd.start,
1064 cmd.end,
1065 Some(cid_resolver),
1066 )
1067 .await?
1068 };
1069
1070 if let Some(instrument_id) = cmd.instrument_id {
1071 reports.retain(|report| report.instrument_id == instrument_id);
1072 }
1073
1074 if cmd.open_only {
1075 reports.retain(|r| r.order_status.is_open());
1076 }
1077
1078 if let Some(start) = cmd.start {
1079 reports.retain(|r| r.ts_last >= start);
1080 }
1081
1082 if let Some(end) = cmd.end {
1083 reports.retain(|r| r.ts_last <= end);
1084 }
1085
1086 Ok(reports)
1087 }
1088
1089 async fn generate_fill_reports(
1090 &self,
1091 cmd: GenerateFillReports,
1092 ) -> anyhow::Result<Vec<FillReport>> {
1093 let mut reports = self
1094 .http_client
1095 .request_fill_reports(self.core.account_id, cmd.start, cmd.end)
1096 .await?;
1097
1098 if let Some(instrument_id) = cmd.instrument_id {
1099 reports.retain(|report| report.instrument_id == instrument_id);
1100 }
1101
1102 if let Some(venue_order_id) = cmd.venue_order_id {
1103 reports.retain(|report| report.venue_order_id.as_str() == venue_order_id.as_str());
1104 }
1105
1106 Ok(reports)
1107 }
1108
1109 async fn generate_position_status_reports(
1110 &self,
1111 cmd: &GeneratePositionStatusReports,
1112 ) -> anyhow::Result<Vec<PositionStatusReport>> {
1113 let mut reports = self
1114 .http_client
1115 .request_position_reports(self.core.account_id)
1116 .await?;
1117
1118 if let Some(instrument_id) = cmd.instrument_id {
1119 reports.retain(|report| report.instrument_id == instrument_id);
1120 }
1121
1122 Ok(reports)
1123 }
1124
1125 async fn generate_mass_status(
1126 &self,
1127 lookback_mins: Option<u64>,
1128 ) -> anyhow::Result<Option<ExecutionMassStatus>> {
1129 log::info!("Generating ExecutionMassStatus (lookback_mins={lookback_mins:?})");
1130
1131 let ts_now = self.clock.get_time_ns();
1132
1133 let start = lookback_mins.map(|mins| {
1134 let lookback_ns = mins * 60 * 1_000_000_000;
1135 UnixNanos::from(ts_now.as_u64().saturating_sub(lookback_ns))
1136 });
1137
1138 let order_cmd = GenerateOrderStatusReports::new(
1139 UUID4::new(),
1140 ts_now,
1141 false, None, start,
1144 None, None, None, );
1148
1149 let fill_cmd = GenerateFillReports::new(
1150 UUID4::new(),
1151 ts_now,
1152 None, None, start,
1155 None, None, None, );
1159
1160 let position_cmd = GeneratePositionStatusReports::new(
1161 UUID4::new(),
1162 ts_now,
1163 None, start,
1165 None, None, None, );
1169
1170 let (order_reports, fill_reports, position_reports) = tokio::try_join!(
1171 self.generate_order_status_reports(&order_cmd),
1172 self.generate_fill_reports(fill_cmd),
1173 self.generate_position_status_reports(&position_cmd),
1174 )?;
1175
1176 log::info!("Received {} OrderStatusReports", order_reports.len());
1177 log::info!("Received {} FillReports", fill_reports.len());
1178 log::info!("Received {} PositionReports", position_reports.len());
1179
1180 let mut mass_status = ExecutionMassStatus::new(
1181 self.core.client_id,
1182 self.core.account_id,
1183 *AX_VENUE,
1184 ts_now,
1185 None,
1186 );
1187
1188 mass_status.add_order_reports(order_reports);
1189 mass_status.add_fill_reports(fill_reports);
1190 mass_status.add_position_reports(position_reports);
1191
1192 Ok(Some(mass_status))
1193 }
1194
1195 fn register_external_order(
1196 &self,
1197 client_order_id: ClientOrderId,
1198 venue_order_id: VenueOrderId,
1199 instrument_id: InstrumentId,
1200 strategy_id: StrategyId,
1201 _ts_init: UnixNanos,
1202 ) {
1203 self.ws_orders.register_external_order(
1204 client_order_id,
1205 venue_order_id,
1206 instrument_id,
1207 strategy_id,
1208 );
1209 }
1210}
1211
1212fn emit_ax_modify_rejected(
1213 emitter: &ExecutionEventEmitter,
1214 clock: &'static AtomicTime,
1215 strategy_id: StrategyId,
1216 instrument_id: InstrumentId,
1217 client_order_id: ClientOrderId,
1218 venue_order_id: Option<VenueOrderId>,
1219 reason: &str,
1220) {
1221 log::warn!("Modify command failed local validation for {client_order_id}: {reason}");
1222 emitter.emit_order_modify_rejected_event(
1223 strategy_id,
1224 instrument_id,
1225 client_order_id,
1226 venue_order_id,
1227 reason,
1228 clock.get_time_ns(),
1229 );
1230}
1231
1232fn dispatch_ws_message(
1234 message: AxOrdersWsMessage,
1235 emitter: &ExecutionEventEmitter,
1236 caches: &OrdersCaches,
1237 account_id: AccountId,
1238 instruments: &AtomicMap<Ustr, InstrumentAny>,
1239 clock: &'static AtomicTime,
1240) {
1241 match message {
1242 AxOrdersWsMessage::Event(event) => {
1243 dispatch_order_event(*event, emitter, caches, account_id, instruments, clock);
1244 }
1245 AxOrdersWsMessage::PlaceOrderResponse(resp) => {
1246 log::debug!(
1247 "Place order response: rid={} oid={}",
1248 resp.rid,
1249 resp.res.oid
1250 );
1251 }
1252 AxOrdersWsMessage::CancelOrderResponse(resp) => {
1253 log::debug!(
1254 "Cancel order response: rid={} accepted={}",
1255 resp.rid,
1256 resp.res.cxl_rx
1257 );
1258 }
1259 AxOrdersWsMessage::OpenOrdersResponse(resp) => {
1260 log::debug!("Open orders response: {} orders", resp.res.orders.len());
1261 }
1262 AxOrdersWsMessage::Error(err) => {
1263 log::warn!("WebSocket error: {}", err.message);
1264 }
1265 AxOrdersWsMessage::Reconnected => {
1266 log::info!("WebSocket reconnected");
1267 }
1268 AxOrdersWsMessage::Authenticated => {
1269 log::debug!("WebSocket authenticated");
1270 }
1271 }
1272}
1273
1274fn dispatch_order_event(
1275 event: AxWsOrderEvent,
1276 emitter: &ExecutionEventEmitter,
1277 caches: &OrdersCaches,
1278 account_id: AccountId,
1279 instruments: &AtomicMap<Ustr, InstrumentAny>,
1280 clock: &'static AtomicTime,
1281) {
1282 match event {
1283 AxWsOrderEvent::Heartbeat => {}
1284 AxWsOrderEvent::Acknowledged(msg) => {
1285 if let Some(event) =
1286 create_order_accepted(&msg.o, msg.ts, msg.tn, caches, account_id, clock)
1287 {
1288 emitter.send_order_event(OrderEventAny::Accepted(event));
1289 } else if let Some(report) = create_order_status_report(
1290 &msg.o,
1291 OrderStatus::Accepted,
1292 msg.ts,
1293 msg.tn,
1294 caches,
1295 account_id,
1296 instruments,
1297 clock,
1298 ) {
1299 emitter.send_order_status_report(report);
1300 }
1301 }
1302 AxWsOrderEvent::PartiallyFilled(msg) => {
1303 dispatch_fill_event(
1304 &msg.o,
1305 &msg.xs,
1306 msg.ts,
1307 msg.tn,
1308 emitter,
1309 caches,
1310 account_id,
1311 instruments,
1312 clock,
1313 );
1314 }
1315 AxWsOrderEvent::Filled(msg) => {
1316 dispatch_fill_event(
1317 &msg.o,
1318 &msg.xs,
1319 msg.ts,
1320 msg.tn,
1321 emitter,
1322 caches,
1323 account_id,
1324 instruments,
1325 clock,
1326 );
1327 cleanup_terminal_order_tracking(&msg.o, caches);
1328 }
1329 AxWsOrderEvent::Canceled(msg) => {
1330 if let Some(event) =
1331 create_order_canceled(&msg.o, msg.ts, msg.tn, caches, account_id, clock)
1332 {
1333 emitter.send_order_event(OrderEventAny::Canceled(event));
1334 } else if let Some(report) = create_order_status_report(
1335 &msg.o,
1336 OrderStatus::Canceled,
1337 msg.ts,
1338 msg.tn,
1339 caches,
1340 account_id,
1341 instruments,
1342 clock,
1343 ) {
1344 emitter.send_order_status_report(report);
1345 }
1346 cleanup_terminal_order_tracking(&msg.o, caches);
1347 }
1348 AxWsOrderEvent::Rejected(msg) => {
1349 let known_reason = msg.r.filter(|r| !matches!(r, AxOrderRejectReason::Unknown));
1350 let reason = known_reason
1351 .as_ref()
1352 .map(AsRef::as_ref)
1353 .or(msg.txt.as_deref())
1354 .unwrap_or("UNKNOWN");
1355
1356 if let Some(event) =
1357 create_order_rejected(&msg.o, reason, msg.ts, msg.tn, caches, account_id, clock)
1358 {
1359 emitter.send_order_event(OrderEventAny::Rejected(event));
1360 }
1361 cleanup_terminal_order_tracking(&msg.o, caches);
1362 }
1363 AxWsOrderEvent::Expired(msg) => {
1364 let as_canceled = matches!(msg.o.tif, AxTimeInForce::Ioc | AxTimeInForce::Fok);
1366 let event = if as_canceled {
1367 create_order_canceled(&msg.o, msg.ts, msg.tn, caches, account_id, clock)
1368 .map(OrderEventAny::Canceled)
1369 } else {
1370 create_order_expired(&msg.o, msg.ts, msg.tn, caches, account_id, clock)
1371 .map(OrderEventAny::Expired)
1372 };
1373
1374 if let Some(event) = event {
1375 emitter.send_order_event(event);
1376 } else if let Some(report) = create_order_status_report(
1377 &msg.o,
1378 if as_canceled {
1379 OrderStatus::Canceled
1380 } else {
1381 OrderStatus::Expired
1382 },
1383 msg.ts,
1384 msg.tn,
1385 caches,
1386 account_id,
1387 instruments,
1388 clock,
1389 ) {
1390 emitter.send_order_status_report(report);
1391 }
1392 cleanup_terminal_order_tracking(&msg.o, caches);
1393 }
1394 AxWsOrderEvent::Replaced(msg) => {
1395 let replacement_venue_order_id = match replacement_venue_order_id(&msg) {
1396 Ok(venue_order_id) => venue_order_id,
1397 Err(e) => {
1398 log::warn!("Invalid AX replace event, awaiting reconciliation: {e}");
1399 return;
1400 }
1401 };
1402
1403 if let Some(event) = create_order_updated(
1404 &msg.no,
1405 &msg.ro,
1406 replacement_venue_order_id,
1407 (msg.ts, msg.tn),
1408 caches,
1409 account_id,
1410 clock,
1411 ) {
1412 emitter.send_order_event(OrderEventAny::Updated(event));
1413 } else if let Some(report) = create_order_status_report(
1414 &msg.no,
1415 OrderStatus::Accepted,
1416 msg.ts,
1417 msg.tn,
1418 caches,
1419 account_id,
1420 instruments,
1421 clock,
1422 ) {
1423 emitter.send_order_status_report(report);
1424 }
1425 }
1426 AxWsOrderEvent::DoneForDay(msg) => {
1427 if let Some(event) =
1428 create_order_expired(&msg.o, msg.ts, msg.tn, caches, account_id, clock)
1429 {
1430 emitter.send_order_event(OrderEventAny::Expired(event));
1431 } else if let Some(report) = create_order_status_report(
1432 &msg.o,
1433 OrderStatus::Expired,
1434 msg.ts,
1435 msg.tn,
1436 caches,
1437 account_id,
1438 instruments,
1439 clock,
1440 ) {
1441 emitter.send_order_status_report(report);
1442 }
1443 cleanup_terminal_order_tracking(&msg.o, caches);
1444 }
1445 AxWsOrderEvent::CancelRejected(msg) => {
1446 let venue_order_id = VenueOrderId::new(&msg.oid);
1447 if let Some(client_order_id) = caches.venue_to_client_id.get(&venue_order_id)
1448 && let Some(metadata) = caches.orders_metadata.get(&client_order_id)
1449 {
1450 let event = OrderCancelRejected::new(
1451 metadata.trader_id,
1452 metadata.strategy_id,
1453 metadata.instrument_id,
1454 metadata.client_order_id,
1455 Ustr::from(msg.r.as_ref()),
1456 UUID4::new(),
1457 clock.get_time_ns(),
1458 metadata.ts_init,
1459 false,
1460 Some(venue_order_id),
1461 Some(account_id),
1462 );
1463 emitter.send_order_event(OrderEventAny::CancelRejected(event));
1464 } else {
1465 log::warn!(
1466 "Could not find metadata for cancel rejected order {}",
1467 msg.oid
1468 );
1469 }
1470 }
1471 }
1472}
1473
1474#[expect(clippy::too_many_arguments)]
1475fn dispatch_fill_event(
1476 order: &AxWsOrder,
1477 execution: &AxWsTradeExecution,
1478 ts: i64,
1479 tn: i64,
1480 emitter: &ExecutionEventEmitter,
1481 caches: &OrdersCaches,
1482 account_id: AccountId,
1483 instruments: &AtomicMap<Ustr, InstrumentAny>,
1484 clock: &'static AtomicTime,
1485) {
1486 if let Some(event) = create_order_filled(order, execution, ts, tn, caches, account_id, clock) {
1487 emitter.send_order_event(OrderEventAny::Filled(event));
1488 } else if let Some(report) = create_fill_report(
1489 order,
1490 execution,
1491 ts,
1492 tn,
1493 caches,
1494 account_id,
1495 instruments,
1496 clock,
1497 ) {
1498 emitter.send_fill_report(report);
1499 }
1500}
1501
1502pub(crate) fn lookup_order_metadata<'a>(
1503 order: &AxWsOrder,
1504 caches: &'a OrdersCaches,
1505) -> Option<dashmap::mapref::one::Ref<'a, ClientOrderId, OrderMetadata>> {
1506 let venue_order_id = VenueOrderId::new(&order.oid);
1507
1508 if let Some(client_order_id) = caches.venue_to_client_id.get(&venue_order_id)
1509 && let Some(metadata) = caches.orders_metadata.get(&*client_order_id)
1510 {
1511 return Some(metadata);
1512 }
1513
1514 if let Some(cid) = order.cid
1515 && let Some(client_order_id) = caches.cid_to_client_order_id.get(&cid)
1516 && let Some(metadata) = caches.orders_metadata.get(&*client_order_id)
1517 {
1518 return Some(metadata);
1519 }
1520
1521 None
1522}
1523
1524pub(crate) fn replacement_venue_order_id(
1525 message: &crate::websocket::messages::AxWsOrderReplaced,
1526) -> anyhow::Result<VenueOrderId> {
1527 if message.noid != message.no.oid {
1528 anyhow::bail!(
1529 "noid '{}' does not match new order oid '{}'",
1530 message.noid,
1531 message.no.oid
1532 );
1533 }
1534
1535 VenueOrderId::new_checked(&message.noid).map_err(anyhow::Error::from)
1536}
1537
1538pub(crate) fn create_order_accepted(
1539 order: &AxWsOrder,
1540 event_ts: i64,
1541 event_tn: i64,
1542 caches: &OrdersCaches,
1543 account_id: AccountId,
1544 clock: &'static AtomicTime,
1545) -> Option<OrderAccepted> {
1546 let venue_order_id = VenueOrderId::new(&order.oid);
1547 let metadata = lookup_order_metadata(order, caches)?;
1548
1549 let client_order_id = metadata.client_order_id;
1550 let trader_id = metadata.trader_id;
1551 let strategy_id = metadata.strategy_id;
1552 let instrument_id = metadata.instrument_id;
1553 drop(metadata);
1554
1555 caches
1556 .venue_to_client_id
1557 .insert(venue_order_id, client_order_id);
1558
1559 if let Some(mut entry) = caches.orders_metadata.get_mut(&client_order_id) {
1560 entry.venue_order_id = Some(venue_order_id);
1561 }
1562
1563 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1564 .map_err(|e| log::error!("{e}"))
1565 .ok()?;
1566
1567 Some(OrderAccepted::new(
1568 trader_id,
1569 strategy_id,
1570 instrument_id,
1571 client_order_id,
1572 venue_order_id,
1573 account_id,
1574 UUID4::new(),
1575 ts_event,
1576 clock.get_time_ns(),
1577 false,
1578 ))
1579}
1580
1581pub(crate) fn create_order_updated(
1582 order: &AxWsOrder,
1583 replaced_order: &AxWsOrder,
1584 replacement_venue_order_id: VenueOrderId,
1585 event_timestamp: (i64, i64),
1586 caches: &OrdersCaches,
1587 account_id: AccountId,
1588 clock: &'static AtomicTime,
1589) -> Option<OrderUpdated> {
1590 let metadata = lookup_order_metadata(order, caches)
1591 .or_else(|| lookup_order_metadata(replaced_order, caches))?;
1592
1593 let client_order_id = metadata.client_order_id;
1594 let trader_id = metadata.trader_id;
1595 let strategy_id = metadata.strategy_id;
1596 let instrument_id = metadata.instrument_id;
1597 let price_precision = metadata.price_precision;
1598 let size_precision = metadata.size_precision;
1599 drop(metadata);
1600
1601 record_replacement_venue_id(caches, client_order_id, replacement_venue_order_id, true);
1602
1603 let ts_event = ax_timestamp_stn_to_unix_nanos(event_timestamp.0, event_timestamp.1)
1604 .map_err(|e| log::error!("{e}"))
1605 .ok()?;
1606
1607 let quantity = Quantity::new(order.q as f64, size_precision);
1608 let price = Price::from_decimal_dp(order.p, price_precision).ok();
1609
1610 Some(OrderUpdated::new(
1611 trader_id,
1612 strategy_id,
1613 instrument_id,
1614 client_order_id,
1615 quantity,
1616 UUID4::new(),
1617 ts_event,
1618 clock.get_time_ns(),
1619 false,
1620 Some(replacement_venue_order_id),
1621 Some(account_id),
1622 price,
1623 None, None, false,
1626 ))
1627}
1628
1629pub(crate) fn create_order_filled(
1630 order: &AxWsOrder,
1631 execution: &AxWsTradeExecution,
1632 event_ts: i64,
1633 event_tn: i64,
1634 caches: &OrdersCaches,
1635 account_id: AccountId,
1636 clock: &'static AtomicTime,
1637) -> Option<OrderFilled> {
1638 let venue_order_id = VenueOrderId::new(&order.oid);
1639 let metadata = lookup_order_metadata(order, caches)?;
1640
1641 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1642 .map_err(|e| log::error!("{e}"))
1643 .ok()?;
1644
1645 let last_qty = Quantity::new(execution.q as f64, metadata.size_precision);
1646 let last_px = Price::from_decimal_dp(execution.p, metadata.price_precision).ok()?;
1647
1648 let order_side = OrderSide::from(order.d);
1649
1650 let liquidity_side = if execution.agg {
1651 LiquiditySide::Taker
1652 } else {
1653 LiquiditySide::Maker
1654 };
1655
1656 Some(OrderFilled::new(
1657 metadata.trader_id,
1658 metadata.strategy_id,
1659 metadata.instrument_id,
1660 metadata.client_order_id,
1661 venue_order_id,
1662 account_id,
1663 TradeId::new(&execution.tid),
1664 order_side,
1665 OrderType::Limit,
1666 last_qty,
1667 last_px,
1668 metadata.quote_currency,
1669 liquidity_side,
1670 UUID4::new(),
1671 ts_event,
1672 clock.get_time_ns(),
1673 false,
1674 None,
1675 None,
1676 None,
1677 ))
1678}
1679
1680pub(crate) fn create_order_canceled(
1681 order: &AxWsOrder,
1682 event_ts: i64,
1683 event_tn: i64,
1684 caches: &OrdersCaches,
1685 account_id: AccountId,
1686 clock: &'static AtomicTime,
1687) -> Option<OrderCanceled> {
1688 let venue_order_id = VenueOrderId::new(&order.oid);
1689 let metadata = lookup_order_metadata(order, caches)?;
1690
1691 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1692 .map_err(|e| log::error!("{e}"))
1693 .ok()?;
1694
1695 Some(OrderCanceled::new(
1696 metadata.trader_id,
1697 metadata.strategy_id,
1698 metadata.instrument_id,
1699 metadata.client_order_id,
1700 UUID4::new(),
1701 ts_event,
1702 clock.get_time_ns(),
1703 false,
1704 Some(venue_order_id),
1705 Some(account_id),
1706 ))
1707}
1708
1709pub(crate) fn create_order_expired(
1710 order: &AxWsOrder,
1711 event_ts: i64,
1712 event_tn: i64,
1713 caches: &OrdersCaches,
1714 account_id: AccountId,
1715 clock: &'static AtomicTime,
1716) -> Option<OrderExpired> {
1717 let venue_order_id = VenueOrderId::new(&order.oid);
1718 let metadata = lookup_order_metadata(order, caches)?;
1719
1720 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1721 .map_err(|e| log::error!("{e}"))
1722 .ok()?;
1723
1724 Some(OrderExpired::new(
1725 metadata.trader_id,
1726 metadata.strategy_id,
1727 metadata.instrument_id,
1728 metadata.client_order_id,
1729 UUID4::new(),
1730 ts_event,
1731 clock.get_time_ns(),
1732 false,
1733 Some(venue_order_id),
1734 Some(account_id),
1735 ))
1736}
1737
1738pub(crate) fn create_order_rejected(
1739 order: &AxWsOrder,
1740 reason: &str,
1741 event_ts: i64,
1742 event_tn: i64,
1743 caches: &OrdersCaches,
1744 account_id: AccountId,
1745 clock: &'static AtomicTime,
1746) -> Option<OrderRejected> {
1747 let metadata = lookup_order_metadata(order, caches)?;
1748
1749 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1750 .map_err(|e| log::error!("{e}"))
1751 .ok()?;
1752 let due_post_only = reason.contains(AX_POST_ONLY_REJECT);
1753
1754 Some(OrderRejected::new(
1755 metadata.trader_id,
1756 metadata.strategy_id,
1757 metadata.instrument_id,
1758 metadata.client_order_id,
1759 account_id,
1760 Ustr::from(reason),
1761 UUID4::new(),
1762 ts_event,
1763 clock.get_time_ns(),
1764 false,
1765 due_post_only,
1766 ))
1767}
1768
1769pub(crate) fn cleanup_terminal_order_tracking(order: &AxWsOrder, caches: &OrdersCaches) {
1770 let venue_order_id = VenueOrderId::new(&order.oid);
1771 let client_order_id = caches
1772 .venue_to_client_id
1773 .remove(&venue_order_id)
1774 .map(|(_, v)| v)
1775 .or_else(|| {
1776 order
1777 .cid
1778 .and_then(|cid| caches.cid_to_client_order_id.remove(&cid).map(|(_, v)| v))
1779 });
1780
1781 if let Some(client_order_id) = client_order_id {
1782 caches.orders_metadata.remove(&client_order_id);
1783 caches
1784 .venue_to_client_id
1785 .retain(|_, mapped_client_order_id| *mapped_client_order_id != client_order_id);
1786 caches
1787 .cid_to_client_order_id
1788 .retain(|_, mapped_client_order_id| *mapped_client_order_id != client_order_id);
1789 }
1790
1791 if let Some(cid) = order.cid {
1792 caches.cid_to_client_order_id.remove(&cid);
1793 }
1794}
1795
1796fn record_replacement_venue_id(
1797 caches: &OrdersCaches,
1798 client_order_id: ClientOrderId,
1799 venue_order_id: VenueOrderId,
1800 remove_previous_venue_ids: bool,
1801) {
1802 caches
1803 .venue_to_client_id
1804 .insert(venue_order_id, client_order_id);
1805 if let Some(mut entry) = caches.orders_metadata.get_mut(&client_order_id) {
1806 entry.venue_order_id = Some(venue_order_id);
1807 }
1808
1809 if remove_previous_venue_ids {
1810 caches
1811 .venue_to_client_id
1812 .retain(|mapped_venue_order_id, mapped_client_order_id| {
1813 *mapped_client_order_id != client_order_id
1814 || *mapped_venue_order_id == venue_order_id
1815 });
1816 }
1817}
1818
1819#[expect(clippy::too_many_arguments)]
1820fn create_order_status_report(
1821 order: &AxWsOrder,
1822 order_status: OrderStatus,
1823 event_ts: i64,
1824 event_tn: i64,
1825 caches: &OrdersCaches,
1826 account_id: AccountId,
1827 instruments: &AtomicMap<Ustr, InstrumentAny>,
1828 clock: &'static AtomicTime,
1829) -> Option<OrderStatusReport> {
1830 let instruments_snap = instruments.load();
1831 let instrument = instruments_snap.get(&order.s)?;
1832 let venue_order_id = VenueOrderId::new(&order.oid);
1833 let instrument_id = instrument.id();
1834 let order_side = OrderSide::from(order.d);
1835 let time_in_force = order.tif.into();
1836
1837 let quantity = Quantity::new(order.q as f64, instrument.size_precision());
1838 let filled_qty = Quantity::new(order.xq as f64, instrument.size_precision());
1839
1840 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1841 .map_err(|e| log::error!("{e}"))
1842 .ok()?;
1843 let ts_init = clock.get_time_ns();
1844
1845 let client_order_id = order.cid.map(|cid| {
1846 caches
1847 .cid_to_client_order_id
1848 .get(&cid)
1849 .map_or_else(|| cid_to_client_order_id(cid), |v| *v)
1850 });
1851
1852 let mut report = OrderStatusReport::new(
1853 account_id,
1854 instrument_id,
1855 client_order_id,
1856 venue_order_id,
1857 order_side.into(),
1858 OrderType::Limit,
1859 time_in_force,
1860 order_status,
1861 quantity,
1862 filled_qty,
1863 ts_event,
1864 ts_event,
1865 ts_init,
1866 Some(UUID4::new()),
1867 );
1868
1869 if let Ok(price) = Price::from_decimal_dp(order.p, instrument.price_precision()) {
1870 report = report.with_price(price);
1871 }
1872
1873 Some(report)
1874}
1875
1876#[expect(clippy::too_many_arguments)]
1877fn create_fill_report(
1878 order: &AxWsOrder,
1879 execution: &AxWsTradeExecution,
1880 event_ts: i64,
1881 event_tn: i64,
1882 caches: &OrdersCaches,
1883 account_id: AccountId,
1884 instruments: &AtomicMap<Ustr, InstrumentAny>,
1885 clock: &'static AtomicTime,
1886) -> Option<FillReport> {
1887 let instruments_snap = instruments.load();
1888 let instrument = instruments_snap.get(&order.s)?;
1889 let venue_order_id = VenueOrderId::new(&order.oid);
1890 let instrument_id = instrument.id();
1891 let order_side = order.d.into();
1892
1893 let last_qty = Quantity::new(execution.q as f64, instrument.size_precision());
1894 let last_px = Price::from_decimal_dp(execution.p, instrument.price_precision()).ok()?;
1895
1896 let liquidity_side = if execution.agg {
1897 LiquiditySide::Taker
1898 } else {
1899 LiquiditySide::Maker
1900 };
1901
1902 let ts_event = ax_timestamp_stn_to_unix_nanos(event_ts, event_tn)
1903 .map_err(|e| log::error!("{e}"))
1904 .ok()?;
1905 let ts_init = clock.get_time_ns();
1906
1907 let client_order_id = order.cid.map(|cid| {
1908 caches
1909 .cid_to_client_order_id
1910 .get(&cid)
1911 .map_or_else(|| cid_to_client_order_id(cid), |v| *v)
1912 });
1913
1914 let commission = Money::zero(instrument.quote_currency());
1918
1919 Some(FillReport::new(
1920 account_id,
1921 instrument_id,
1922 venue_order_id,
1923 TradeId::new(&execution.tid),
1924 order_side,
1925 last_qty,
1926 last_px,
1927 commission,
1928 liquidity_side,
1929 client_order_id,
1930 None,
1931 ts_event,
1932 ts_init,
1933 Some(UUID4::new()),
1934 ))
1935}
1936
1937fn validate_order_for_ax_submit(order: &OrderAny) -> anyhow::Result<()> {
1938 if !matches!(order.order_type(), OrderType::Market | OrderType::Limit) {
1939 anyhow::bail!(
1940 "Unsupported order type: {:?}, the Architect AX adapter accepts Nautilus MARKET and LIMIT orders",
1941 order.order_type(),
1942 );
1943 }
1944
1945 if !matches!(
1947 order.time_in_force(),
1948 TimeInForce::Gtc | TimeInForce::Ioc | TimeInForce::Day
1949 ) {
1950 anyhow::bail!(
1951 "Unsupported time in force: {:?}, AX supports GTC, IOC, and DAY",
1952 order.time_in_force(),
1953 );
1954 }
1955
1956 validate_order_instructions(
1957 order.is_reduce_only(),
1958 order.is_quote_quantity(),
1959 order.display_qty().is_some(),
1960 )?;
1961
1962 quantity_to_contracts(order.quantity())?;
1963
1964 Ok(())
1965}
1966
1967fn validate_order_init_instructions(order_init: &OrderInitialized) -> anyhow::Result<()> {
1968 validate_order_instructions(
1969 order_init.reduce_only,
1970 order_init.quote_quantity,
1971 order_init.display_qty.is_some(),
1972 )
1973}
1974
1975fn validate_order_instructions(
1976 reduce_only: bool,
1977 quote_quantity: bool,
1978 has_display_qty: bool,
1979) -> anyhow::Result<()> {
1980 if reduce_only {
1981 anyhow::bail!("AX does not support reduce-only orders");
1982 }
1983
1984 if quote_quantity {
1985 anyhow::bail!(
1986 "Architect AX adapter cannot encode quote_quantity; submit a base quantity instead"
1987 );
1988 }
1989
1990 if has_display_qty {
1991 anyhow::bail!("Architect AX adapter cannot encode display_qty iceberg instructions");
1992 }
1993
1994 Ok(())
1995}
1996
1997fn classify_ax_http_failure(error: &AxHttpError) -> CommandFailure {
2002 let message = error.to_string();
2003 match error {
2004 AxHttpError::MissingCredentials
2005 | AxHttpError::MissingSessionToken
2006 | AxHttpError::ValidationError(_)
2007 | AxHttpError::BuildError(_) => CommandFailure::NotSent(message),
2008 AxHttpError::ApiError { .. }
2009 | AxHttpError::JsonError(_)
2010 | AxHttpError::Canceled(_)
2011 | AxHttpError::NetworkError(_)
2012 | AxHttpError::UnexpectedStatus { .. } => CommandFailure::Ambiguous(message),
2013 }
2014}
2015
2016fn classify_ax_ws_failure(error: &AxOrdersWsClientError) -> CommandFailure {
2017 match error {
2018 AxOrdersWsClientError::ClientError(message) => CommandFailure::NotSent(message.clone()),
2019 AxOrdersWsClientError::Transport(_)
2020 | AxOrdersWsClientError::ChannelError(_)
2021 | AxOrdersWsClientError::AuthenticationError(_) => {
2022 CommandFailure::Ambiguous(error.to_string())
2023 }
2024 }
2025}
2026
2027#[cfg(test)]
2028mod tests {
2029 use std::sync::Arc;
2030
2031 use dashmap::DashMap;
2032 use nautilus_common::messages::ExecutionEvent;
2033 use nautilus_core::time::get_atomic_clock_realtime;
2034 use nautilus_model::{
2035 identifiers::{AccountId, ClientOrderId, InstrumentId, StrategyId, TraderId, VenueOrderId},
2036 orders::builder::OrderTestBuilder,
2037 types::{Currency, Price, Quantity},
2038 };
2039 use rstest::rstest;
2040 use rust_decimal::Decimal;
2041 use rust_decimal_macros::dec;
2042 use ustr::Ustr;
2043
2044 use super::*;
2045 use crate::{
2046 common::enums::{AxOrderSide, AxOrderStatus, AxTimeInForce},
2047 http::error::AxBuildError,
2048 websocket::{
2049 messages::{AxWsOrderExpired, AxWsTradeExecution, OrderMetadata},
2050 orders::OrdersCaches,
2051 },
2052 };
2053
2054 fn test_caches() -> OrdersCaches {
2055 OrdersCaches {
2056 orders_metadata: Arc::new(DashMap::new()),
2057 venue_to_client_id: Arc::new(DashMap::new()),
2058 cid_to_client_order_id: Arc::new(DashMap::new()),
2059 }
2060 }
2061
2062 fn test_ws_order(oid: &str, price: Decimal, qty: u64) -> AxWsOrder {
2063 AxWsOrder {
2064 oid: oid.to_string(),
2065 u: "user".to_string(),
2066 s: Ustr::from("BTC-PERP"),
2067 p: price,
2068 q: qty,
2069 xq: 0,
2070 rq: qty,
2071 o: AxOrderStatus::Accepted,
2072 d: AxOrderSide::Buy,
2073 tif: AxTimeInForce::Gtc,
2074 ts: 1609459200,
2075 tn: 0,
2076 cid: None,
2077 tag: None,
2078 txt: None,
2079 }
2080 }
2081
2082 #[rstest]
2083 fn test_create_order_updated_uses_ws_replacement_id_before_http_response() {
2084 let caches = test_caches();
2085 let clock = get_atomic_clock_realtime();
2086 let account_id = AccountId::from("AX-001");
2087 let client_order_id = ClientOrderId::from("O-WS-FIRST");
2088 let old_venue_order_id = VenueOrderId::new("OLD-OID");
2089 let new_venue_order_id = VenueOrderId::new("NEW-OID");
2090
2091 let mut metadata = test_metadata(client_order_id, InstrumentId::from("BTC-PERP.AX"));
2092 metadata.venue_order_id = Some(old_venue_order_id);
2093 caches.orders_metadata.insert(client_order_id, metadata);
2094 caches
2095 .venue_to_client_id
2096 .insert(old_venue_order_id, client_order_id);
2097
2098 let old_order = test_ws_order(old_venue_order_id.as_str(), dec!(50000.00), 100);
2099 let new_order = test_ws_order(new_venue_order_id.as_str(), dec!(50001.00), 100);
2100 let event = create_order_updated(
2101 &new_order,
2102 &old_order,
2103 new_venue_order_id,
2104 (1609459200, 0),
2105 &caches,
2106 account_id,
2107 clock,
2108 )
2109 .expect("should produce OrderUpdated");
2110
2111 assert_eq!(event.venue_order_id, Some(new_venue_order_id));
2112 assert_eq!(
2113 caches
2114 .orders_metadata
2115 .get(&client_order_id)
2116 .unwrap()
2117 .venue_order_id,
2118 Some(new_venue_order_id),
2119 );
2120 assert!(!caches.venue_to_client_id.contains_key(&old_venue_order_id));
2121 assert_eq!(
2122 *caches.venue_to_client_id.get(&new_venue_order_id).unwrap(),
2123 client_order_id,
2124 );
2125 }
2126
2127 #[rstest]
2128 fn test_record_http_replacement_retains_previous_venue_id_until_ws_event() {
2129 let caches = test_caches();
2130 let client_order_id = ClientOrderId::from("O-HTTP-FIRST");
2131 let old_venue_order_id = VenueOrderId::new("OLD-OID");
2132 let new_venue_order_id = VenueOrderId::new("NEW-OID");
2133 let mut metadata = test_metadata(client_order_id, InstrumentId::from("BTC-PERP.AX"));
2134 metadata.venue_order_id = Some(old_venue_order_id);
2135 caches.orders_metadata.insert(client_order_id, metadata);
2136 caches
2137 .venue_to_client_id
2138 .insert(old_venue_order_id, client_order_id);
2139
2140 record_replacement_venue_id(&caches, client_order_id, new_venue_order_id, false);
2141
2142 assert!(caches.venue_to_client_id.contains_key(&old_venue_order_id));
2143 assert!(caches.venue_to_client_id.contains_key(&new_venue_order_id));
2144 assert_eq!(
2145 caches
2146 .orders_metadata
2147 .get(&client_order_id)
2148 .unwrap()
2149 .venue_order_id,
2150 Some(new_venue_order_id),
2151 );
2152
2153 record_replacement_venue_id(&caches, client_order_id, new_venue_order_id, true);
2154
2155 assert!(!caches.venue_to_client_id.contains_key(&old_venue_order_id));
2156 assert!(caches.venue_to_client_id.contains_key(&new_venue_order_id));
2157 }
2158
2159 fn test_metadata(client_order_id: ClientOrderId, instrument_id: InstrumentId) -> OrderMetadata {
2160 OrderMetadata {
2161 trader_id: TraderId::from("TRADER-001"),
2162 strategy_id: StrategyId::from("S-001"),
2163 instrument_id,
2164 client_order_id,
2165 venue_order_id: None,
2166 ts_init: 0.into(),
2167 size_precision: 0,
2168 price_precision: 2,
2169 quote_currency: Currency::USD(),
2170 }
2171 }
2172
2173 fn test_execution(tid: &str, price: Decimal, qty: u64, agg: bool) -> AxWsTradeExecution {
2174 AxWsTradeExecution {
2175 tid: tid.to_string(),
2176 s: Ustr::from("BTC-PERP"),
2177 q: qty,
2178 p: price,
2179 d: AxOrderSide::Buy,
2180 agg,
2181 }
2182 }
2183
2184 #[rstest]
2185 fn test_create_order_accepted_populates_cache_and_event() {
2186 let caches = test_caches();
2187 let clock = get_atomic_clock_realtime();
2188 let account_id = AccountId::from("AX-001");
2189 let client_order_id = ClientOrderId::from("O-ACK");
2190 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2191 let venue_order_id = VenueOrderId::new("OID-ACK");
2192
2193 caches.orders_metadata.insert(
2194 client_order_id,
2195 test_metadata(client_order_id, instrument_id),
2196 );
2197 let cid_value = 7u64;
2198 caches
2199 .cid_to_client_order_id
2200 .insert(cid_value, client_order_id);
2201
2202 let mut ws_order = test_ws_order(venue_order_id.as_str(), dec!(50500.00), 100);
2203 ws_order.cid = Some(cid_value);
2204
2205 let event = create_order_accepted(&ws_order, 1609459200, 500, &caches, account_id, clock)
2206 .expect("should produce OrderAccepted");
2207
2208 assert_eq!(event.venue_order_id, venue_order_id);
2209 assert_eq!(event.client_order_id, client_order_id);
2210 assert_eq!(event.account_id, account_id);
2211 assert_eq!(event.instrument_id, instrument_id);
2212 assert_eq!(event.trader_id, TraderId::from("TRADER-001"));
2213 assert_eq!(event.strategy_id, StrategyId::from("S-001"));
2214 assert_eq!(
2215 event.ts_event,
2216 UnixNanos::from(1_609_459_200_000_000_500u64)
2217 );
2218
2219 assert_eq!(
2221 *caches.venue_to_client_id.get(&venue_order_id).unwrap(),
2222 client_order_id,
2223 );
2224 let meta = caches.orders_metadata.get(&client_order_id).unwrap();
2225 assert_eq!(meta.venue_order_id, Some(venue_order_id));
2226 }
2227
2228 #[rstest]
2229 fn test_create_order_accepted_returns_none_without_metadata() {
2230 let caches = test_caches();
2231 let clock = get_atomic_clock_realtime();
2232 let account_id = AccountId::from("AX-001");
2233 let ws_order = test_ws_order("OID-UNKNOWN", dec!(100.00), 10);
2234
2235 let result = create_order_accepted(&ws_order, 1609459200, 0, &caches, account_id, clock);
2236 assert!(result.is_none());
2237 assert!(caches.venue_to_client_id.is_empty());
2238 }
2239
2240 #[rstest]
2241 fn test_lookup_order_metadata_cid_fallback() {
2242 let caches = test_caches();
2243 let client_order_id = ClientOrderId::from("O-CID");
2244 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2245 caches.orders_metadata.insert(
2246 client_order_id,
2247 test_metadata(client_order_id, instrument_id),
2248 );
2249 caches.cid_to_client_order_id.insert(99, client_order_id);
2250
2251 let mut ws_order = test_ws_order("UNKNOWN-OID", dec!(0), 0);
2252 ws_order.cid = Some(99);
2253
2254 let found = lookup_order_metadata(&ws_order, &caches).expect("cid fallback should find");
2255 assert_eq!(found.client_order_id, client_order_id);
2256 }
2257
2258 #[rstest]
2259 fn test_lookup_order_metadata_returns_none_when_unknown() {
2260 let caches = test_caches();
2261 let ws_order = test_ws_order("UNKNOWN-OID", dec!(0), 0);
2262 assert!(lookup_order_metadata(&ws_order, &caches).is_none());
2263 }
2264
2265 #[rstest]
2266 #[case(true, LiquiditySide::Taker)]
2267 #[case(false, LiquiditySide::Maker)]
2268 fn test_create_order_filled_maps_liquidity_side(
2269 #[case] agg: bool,
2270 #[case] expected: LiquiditySide,
2271 ) {
2272 let caches = test_caches();
2273 let clock = get_atomic_clock_realtime();
2274 let account_id = AccountId::from("AX-001");
2275 let client_order_id = ClientOrderId::from("O-FILL");
2276 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2277 let venue_order_id = VenueOrderId::new("OID-FILL");
2278
2279 caches.orders_metadata.insert(
2280 client_order_id,
2281 test_metadata(client_order_id, instrument_id),
2282 );
2283 caches
2284 .venue_to_client_id
2285 .insert(venue_order_id, client_order_id);
2286
2287 let order = test_ws_order(venue_order_id.as_str(), dec!(50500.00), 100);
2288 let execution = test_execution("TID-1", dec!(50500.00), 25, agg);
2289
2290 let event = create_order_filled(
2291 &order, &execution, 1609459200, 0, &caches, account_id, clock,
2292 )
2293 .expect("should produce OrderFilled");
2294
2295 assert_eq!(event.venue_order_id, venue_order_id);
2296 assert_eq!(event.client_order_id, client_order_id);
2297 assert_eq!(event.trade_id, TradeId::new("TID-1"));
2298 assert_eq!(event.last_qty, Quantity::new(25.0, 0));
2299 assert_eq!(event.last_px, Price::from("50500.00"));
2300 assert_eq!(event.liquidity_side, expected);
2301 }
2302
2303 #[rstest]
2304 fn test_create_order_canceled_populates_identifiers() {
2305 let caches = test_caches();
2306 let clock = get_atomic_clock_realtime();
2307 let account_id = AccountId::from("AX-001");
2308 let client_order_id = ClientOrderId::from("O-CXL");
2309 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2310 let venue_order_id = VenueOrderId::new("OID-CXL");
2311
2312 caches.orders_metadata.insert(
2313 client_order_id,
2314 test_metadata(client_order_id, instrument_id),
2315 );
2316 caches
2317 .venue_to_client_id
2318 .insert(venue_order_id, client_order_id);
2319
2320 let order = test_ws_order(venue_order_id.as_str(), dec!(100.00), 10);
2321 let event = create_order_canceled(&order, 1609459200, 0, &caches, account_id, clock)
2322 .expect("should produce OrderCanceled");
2323
2324 assert_eq!(event.venue_order_id, Some(venue_order_id));
2325 assert_eq!(event.client_order_id, client_order_id);
2326 assert_eq!(event.account_id, Some(account_id));
2327 assert_eq!(event.instrument_id, instrument_id);
2328 }
2329
2330 #[rstest]
2331 fn test_create_order_expired_populates_identifiers() {
2332 let caches = test_caches();
2333 let clock = get_atomic_clock_realtime();
2334 let account_id = AccountId::from("AX-001");
2335 let client_order_id = ClientOrderId::from("O-EXP");
2336 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2337 let venue_order_id = VenueOrderId::new("OID-EXP");
2338
2339 caches.orders_metadata.insert(
2340 client_order_id,
2341 test_metadata(client_order_id, instrument_id),
2342 );
2343 caches
2344 .venue_to_client_id
2345 .insert(venue_order_id, client_order_id);
2346
2347 let order = test_ws_order(venue_order_id.as_str(), dec!(100.00), 10);
2348 let event = create_order_expired(&order, 1609459200, 0, &caches, account_id, clock)
2349 .expect("should produce OrderExpired");
2350
2351 assert_eq!(event.venue_order_id, Some(venue_order_id));
2352 assert_eq!(event.client_order_id, client_order_id);
2353 }
2354
2355 #[rstest]
2356 #[case(AxTimeInForce::Ioc, true)]
2357 #[case(AxTimeInForce::Fok, true)]
2358 #[case(AxTimeInForce::Day, false)]
2359 #[case(AxTimeInForce::Gtc, false)]
2360 fn test_dispatch_expired_maps_ioc_fok_to_canceled(
2361 #[case] tif: AxTimeInForce,
2362 #[case] expect_canceled: bool,
2363 ) {
2364 let clock = get_atomic_clock_realtime();
2365 let account_id = AccountId::from("AX-001");
2366 let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
2367 let mut emitter = ExecutionEventEmitter::new(
2368 clock,
2369 TraderId::from("TESTER-001"),
2370 account_id,
2371 AccountType::Margin,
2372 None,
2373 );
2374 emitter.set_sender(tx);
2375
2376 let caches = test_caches();
2377 let client_order_id = ClientOrderId::from("O-EXP");
2378 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2379 caches.orders_metadata.insert(
2380 client_order_id,
2381 test_metadata(client_order_id, instrument_id),
2382 );
2383 caches
2384 .venue_to_client_id
2385 .insert(VenueOrderId::new("OID-EXP"), client_order_id);
2386
2387 let mut order = test_ws_order("OID-EXP", dec!(100.00), 10);
2388 order.tif = tif;
2389 let event = AxWsOrderEvent::Expired(AxWsOrderExpired {
2390 ts: 1609459200,
2391 tn: 0,
2392 eid: "E-EXP".to_string(),
2393 o: order,
2394 });
2395
2396 let instruments: AtomicMap<Ustr, InstrumentAny> = AtomicMap::new();
2397 dispatch_order_event(event, &emitter, &caches, account_id, &instruments, clock);
2398
2399 match rx.try_recv().expect("an order event should be emitted") {
2400 ExecutionEvent::Order(OrderEventAny::Canceled(_)) => assert!(expect_canceled),
2401 ExecutionEvent::Order(OrderEventAny::Expired(_)) => assert!(!expect_canceled),
2402 other => panic!("unexpected event: {other:?}"),
2403 }
2404 }
2405
2406 #[rstest]
2407 fn test_create_order_rejected_sets_due_post_only_when_reason_matches() {
2408 let caches = test_caches();
2409 let clock = get_atomic_clock_realtime();
2410 let account_id = AccountId::from("AX-001");
2411 let client_order_id = ClientOrderId::from("O-REJ");
2412 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2413
2414 caches.orders_metadata.insert(
2415 client_order_id,
2416 test_metadata(client_order_id, instrument_id),
2417 );
2418 caches
2419 .venue_to_client_id
2420 .insert(VenueOrderId::new("OID-REJ"), client_order_id);
2421
2422 let order = test_ws_order("OID-REJ", dec!(100.00), 10);
2423 let reason = "post-only order would cross the book";
2425 let event =
2426 create_order_rejected(&order, reason, 1609459200, 0, &caches, account_id, clock)
2427 .expect("should produce OrderRejected");
2428
2429 assert!(event.due_post_only, "post-only reason should set flag");
2430 assert_eq!(event.reason, Ustr::from(reason));
2431 }
2432
2433 #[rstest]
2434 fn test_create_order_rejected_clears_due_post_only_for_other_reasons() {
2435 let caches = test_caches();
2436 let clock = get_atomic_clock_realtime();
2437 let account_id = AccountId::from("AX-001");
2438 let client_order_id = ClientOrderId::from("O-REJ-2");
2439 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2440
2441 caches.orders_metadata.insert(
2442 client_order_id,
2443 test_metadata(client_order_id, instrument_id),
2444 );
2445 caches
2446 .venue_to_client_id
2447 .insert(VenueOrderId::new("OID-REJ-2"), client_order_id);
2448
2449 let order = test_ws_order("OID-REJ-2", dec!(100.00), 10);
2450 let event = create_order_rejected(
2451 &order,
2452 "INSUFFICIENT_MARGIN",
2453 1609459200,
2454 0,
2455 &caches,
2456 account_id,
2457 clock,
2458 )
2459 .expect("should produce OrderRejected");
2460
2461 assert!(!event.due_post_only);
2462 assert_eq!(event.reason, Ustr::from("INSUFFICIENT_MARGIN"));
2463 }
2464
2465 #[rstest]
2466 fn test_cleanup_terminal_order_tracking_removes_all_caches() {
2467 let caches = test_caches();
2468 let client_order_id = ClientOrderId::from("O-CLEAN");
2469 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2470 let venue_order_id = VenueOrderId::new("OID-CLEAN");
2471 let stale_venue_order_id = VenueOrderId::new("OID-CLEAN-OLD");
2472 let cid_value = 123u64;
2473
2474 caches.orders_metadata.insert(
2475 client_order_id,
2476 test_metadata(client_order_id, instrument_id),
2477 );
2478 caches
2479 .venue_to_client_id
2480 .insert(venue_order_id, client_order_id);
2481 caches
2482 .venue_to_client_id
2483 .insert(stale_venue_order_id, client_order_id);
2484 caches
2485 .cid_to_client_order_id
2486 .insert(cid_value, client_order_id);
2487
2488 let mut order = test_ws_order(venue_order_id.as_str(), dec!(100.00), 10);
2489 order.cid = Some(cid_value);
2490
2491 cleanup_terminal_order_tracking(&order, &caches);
2492
2493 assert!(caches.orders_metadata.is_empty());
2494 assert!(caches.venue_to_client_id.is_empty());
2495 assert!(caches.cid_to_client_order_id.is_empty());
2496 }
2497
2498 #[rstest]
2499 fn test_cleanup_terminal_order_tracking_via_cid_when_venue_missing() {
2500 let caches = test_caches();
2501 let client_order_id = ClientOrderId::from("O-CLEAN-CID");
2502 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2503 let cid_value = 321u64;
2504
2505 caches.orders_metadata.insert(
2506 client_order_id,
2507 test_metadata(client_order_id, instrument_id),
2508 );
2509 caches
2510 .cid_to_client_order_id
2511 .insert(cid_value, client_order_id);
2512
2513 let mut order = test_ws_order("OID-UNKNOWN", dec!(100.00), 10);
2515 order.cid = Some(cid_value);
2516
2517 cleanup_terminal_order_tracking(&order, &caches);
2518
2519 assert!(caches.orders_metadata.is_empty());
2520 assert!(caches.cid_to_client_order_id.is_empty());
2521 }
2522
2523 #[rstest]
2524 fn test_cleanup_terminal_order_tracking_noop_when_unknown() {
2525 let caches = test_caches();
2526 let other = ClientOrderId::from("OTHER");
2527 let instrument_id = InstrumentId::from("BTC-PERP.AX");
2528 caches
2529 .orders_metadata
2530 .insert(other, test_metadata(other, instrument_id));
2531
2532 let order = test_ws_order("OID-NOT-TRACKED", dec!(100.00), 10);
2533 cleanup_terminal_order_tracking(&order, &caches);
2534
2535 assert_eq!(caches.orders_metadata.len(), 1);
2537 }
2538
2539 #[rstest]
2540 fn test_cancel_on_disconnect_url_no_existing_query() {
2541 let mut url = "wss://example.com/orders/ws".to_string();
2542 let separator = if url.contains('?') { "&" } else { "?" };
2543 url.push_str(&format!("{separator}cancel_on_disconnect=true"));
2544 assert_eq!(url, "wss://example.com/orders/ws?cancel_on_disconnect=true");
2545 }
2546
2547 #[rstest]
2548 fn test_cancel_on_disconnect_url_with_existing_query() {
2549 let mut url = "wss://example.com/orders/ws?token=abc".to_string();
2550 let separator = if url.contains('?') { "&" } else { "?" };
2551 url.push_str(&format!("{separator}cancel_on_disconnect=true"));
2552 assert_eq!(
2553 url,
2554 "wss://example.com/orders/ws?token=abc&cancel_on_disconnect=true"
2555 );
2556 }
2557
2558 fn limit_order_for_validation(quantity: Quantity) -> OrderAny {
2559 OrderTestBuilder::new(OrderType::Limit)
2560 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2561 .side(OrderSide::Buy)
2562 .quantity(quantity)
2563 .price(Price::from("1.10"))
2564 .build()
2565 }
2566
2567 #[rstest]
2568 fn test_validate_order_for_ax_submit_accepts_supported_limit_order() {
2569 let order = limit_order_for_validation(Quantity::from("10"));
2570
2571 assert!(validate_order_for_ax_submit(&order).is_ok());
2572 }
2573
2574 #[rstest]
2575 fn test_validate_order_for_ax_submit_denies_unsupported_order_type() {
2576 let order = OrderTestBuilder::new(OrderType::StopMarket)
2577 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2578 .side(OrderSide::Buy)
2579 .quantity(Quantity::from("10"))
2580 .trigger_price(Price::from("1.10"))
2581 .build();
2582
2583 let err = validate_order_for_ax_submit(&order).unwrap_err();
2584
2585 assert!(err.to_string().contains("Unsupported order type"));
2586 }
2587
2588 #[rstest]
2589 fn test_validate_order_for_ax_submit_denies_gtd_time_in_force() {
2590 let order = OrderTestBuilder::new(OrderType::Limit)
2591 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2592 .side(OrderSide::Buy)
2593 .quantity(Quantity::from("10"))
2594 .price(Price::from("1.10"))
2595 .time_in_force(TimeInForce::Gtd)
2596 .expire_time(UnixNanos::from(2_000_000_000_000_000_000u64))
2597 .build();
2598
2599 let err = validate_order_for_ax_submit(&order).unwrap_err();
2600
2601 assert!(err.to_string().contains("Unsupported time in force"));
2602 }
2603
2604 #[rstest]
2605 fn test_validate_order_for_ax_submit_denies_stop_limit_order() {
2606 let order = OrderTestBuilder::new(OrderType::StopLimit)
2607 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2608 .side(OrderSide::Buy)
2609 .quantity(Quantity::from("10"))
2610 .price(Price::from("1.11"))
2611 .trigger_price(Price::from("1.10"))
2612 .build();
2613
2614 let err = validate_order_for_ax_submit(&order).unwrap_err();
2615
2616 assert!(err.to_string().contains("Unsupported order type"));
2617 }
2618
2619 #[rstest]
2620 fn test_validate_order_for_ax_submit_denies_fok_time_in_force() {
2621 let order = OrderTestBuilder::new(OrderType::Limit)
2622 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2623 .side(OrderSide::Buy)
2624 .quantity(Quantity::from("10"))
2625 .price(Price::from("1.10"))
2626 .time_in_force(TimeInForce::Fok)
2627 .build();
2628
2629 let err = validate_order_for_ax_submit(&order).unwrap_err();
2630
2631 assert!(err.to_string().contains("Unsupported time in force"));
2632 }
2633
2634 #[rstest]
2635 fn test_validate_order_for_ax_submit_denies_fractional_quantity() {
2636 let order = limit_order_for_validation(Quantity::from("10.5"));
2637
2638 let err = validate_order_for_ax_submit(&order).unwrap_err();
2639
2640 assert!(err.to_string().contains("whole contract"));
2641 }
2642
2643 #[rstest]
2644 fn test_validate_order_for_ax_submit_denies_quote_quantity() {
2645 let order = OrderTestBuilder::new(OrderType::Market)
2646 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2647 .side(OrderSide::Buy)
2648 .quantity(Quantity::from("10"))
2649 .time_in_force(TimeInForce::Ioc)
2650 .quote_quantity(true)
2651 .build();
2652
2653 let err = validate_order_for_ax_submit(&order).unwrap_err();
2654
2655 assert!(err.to_string().contains("quote_quantity"));
2656 }
2657
2658 #[rstest]
2659 fn test_validate_order_for_ax_submit_denies_display_quantity() {
2660 let order = OrderTestBuilder::new(OrderType::Limit)
2661 .instrument_id(InstrumentId::from("EURUSD-PERP.AX"))
2662 .side(OrderSide::Buy)
2663 .quantity(Quantity::from("10"))
2664 .price(Price::from("1.10"))
2665 .display_qty(Quantity::from("5"))
2666 .build();
2667
2668 let err = validate_order_for_ax_submit(&order).unwrap_err();
2669
2670 assert!(err.to_string().contains("display_qty"));
2671 }
2672
2673 #[rstest]
2674 #[case(AxHttpError::MissingCredentials, true)]
2675 #[case(AxHttpError::MissingSessionToken, true)]
2676 #[case(AxHttpError::ValidationError("bad param".to_string()), true)]
2677 #[case(AxHttpError::BuildError(AxBuildError::MissingOrderId), true)]
2678 #[case(AxHttpError::ApiError { message: "invalid modification".to_string() }, false)]
2679 #[case(AxHttpError::JsonError("parse failure".to_string()), false)]
2680 #[case(AxHttpError::Canceled("shutdown".to_string()), false)]
2681 #[case(AxHttpError::NetworkError("timeout".to_string()), false)]
2682 #[case(AxHttpError::UnexpectedStatus { status: 400, body: "invalid".to_string() }, false)]
2683 #[case(AxHttpError::UnexpectedStatus { status: 503, body: String::new() }, false)]
2684 fn test_classify_ax_http_failure(#[case] error: AxHttpError, #[case] expect_not_sent: bool) {
2685 let expected = if expect_not_sent {
2687 CommandFailure::NotSent(error.to_string())
2688 } else {
2689 CommandFailure::Ambiguous(error.to_string())
2690 };
2691
2692 assert_eq!(classify_ax_http_failure(&error), expected);
2693 }
2694
2695 #[rstest]
2696 #[case(
2697 AxOrdersWsClientError::ClientError("missing venue_order_id".to_string()),
2698 Some("missing venue_order_id")
2699 )]
2700 #[case(AxOrdersWsClientError::ChannelError("handler closed".to_string()), None)]
2701 #[case(AxOrdersWsClientError::Transport("connection reset".to_string()), None)]
2702 #[case(AxOrdersWsClientError::AuthenticationError("token expired".to_string()), None)]
2703 fn test_classify_ax_ws_failure(
2704 #[case] error: AxOrdersWsClientError,
2705 #[case] not_sent_reason: Option<&str>,
2706 ) {
2707 let expected = match not_sent_reason {
2709 Some(reason) => CommandFailure::NotSent(reason.to_string()),
2710 None => CommandFailure::Ambiguous(error.to_string()),
2711 };
2712
2713 assert_eq!(classify_ax_ws_failure(&error), expected);
2714 }
2715}