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