Skip to main content

nautilus_architect_ax/
execution.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Live execution client implementation for the AX Exchange adapter.
17
18use 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/// Live execution client for the AX Exchange.
90#[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    /// Creates a new [`AxExecutionClient`].
105    ///
106    /// # Errors
107    ///
108    /// Returns an error if the client fails to initialize.
109    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            // AX emulates market orders with preview-priced IOC limits, so book moves
248            // between preview and submission can produce partial fills.
249            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                    // AX classifies no send failure as a venue rejection (see
341                    // `classify_ax_ws_failure`); both terminal-valid classes reject the same way.
342                    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        // `OrderCancelRejected` comes only from the WS CancelRejected event;
371        // local send failures leave the outcome to reconciliation.
372        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    /// Polls the cache until the account is registered or timeout is reached.
461    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        // Reset so requests work after a previous disconnect
540        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        // Instruments load after authenticating because their fee rates come from the
553        // authenticated `/whoami`. A zero-fee fallback would outlive the failure that caused it,
554        // since `set_instruments_initialized` stops a reconnect from retrying the load.
555        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        // Read immutable order fields from cache before spawning
688        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                // No replace failure is an unambiguous rejection (see
934                // `classify_ax_http_failure`); leave the order pending.
935                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        // Snapshot open orders so we can emit cancel events after the HTTP request
975        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                    // AX does not push WS cancel confirmations for HTTP-initiated
992                    // cancels, so emit OrderCanceled events locally and clean up
993                    // tracking state to prevent duplicates if WS events arrive
994                    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                // A whole-request failure has no per-order venue results and
1024                // must not fan out per-order rejections.
1025                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, // open_only
1170            None,  // instrument_id
1171            start,
1172            None, // end
1173            None, // params
1174            None, // correlation_id
1175        );
1176
1177        let fill_cmd = GenerateFillReports::new(
1178            UUID4::new(),
1179            ts_now,
1180            None, // instrument_id
1181            None, // venue_order_id
1182            start,
1183            None, // end
1184            None, // params
1185            None, // correlation_id
1186        );
1187
1188        let position_cmd = GeneratePositionStatusReports::new(
1189            UUID4::new(),
1190            ts_now,
1191            None, // instrument_id
1192            start,
1193            None, // end
1194            None, // params
1195            None, // correlation_id
1196        );
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
1260/// Dispatches a WebSocket message using the event emitter.
1261fn 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            // AX reports an unfilled IOC/FOK as EXPIRED; Nautilus models those as canceled
1393            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, // trigger_price
1652        None, // protection_price
1653        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    // The WS trade execution payload does not include fee data so
1966    // commission is zero here. The REST /fills endpoint (used during
1967    // reconciliation via parse_fill_report) includes accurate fees.
1968    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    // AX accepts only GTC, IOC, and DAY; deny others locally to avoid an opaque venue error
1997    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
2048// The AX HTTP API documents only a bare 400 with no error schema, so no venue
2049// failure can be allowlisted as an unambiguous rejection; those arrive via the
2050// orders WS `Rejected` and `CancelRejected` events instead. AX therefore never
2051// classifies a send failure as `CommandFailure::VenueRejected`.
2052fn 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        // Side effects on caches
2375        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        // Use the literal venue reason text so this guards the AX_POST_ONLY_REJECT constant
2579        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        // Venue id missing from cache
2669        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        // Unrelated metadata still present
2691        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        // Asserting the exact variant also pins that AX never classifies a venue rejection
3158        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        // Asserting the exact variant also pins that AX never classifies a venue rejection
3180        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}