Skip to main content

nautilus_architect_ax/
data.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 market data client implementation for the AX Exchange adapter.
17
18use std::{
19    future::Future,
20    sync::{
21        Arc,
22        atomic::{AtomicBool, Ordering},
23    },
24    time::Duration,
25};
26
27use ahash::{AHashMap, AHashSet};
28use anyhow::Context;
29use async_trait::async_trait;
30use futures_util::StreamExt;
31use jiff::{SignedDuration, Timestamp};
32use nautilus_common::{
33    clients::DataClient,
34    live::{runner::get_data_event_sender, sender::EventSender},
35    messages::{
36        DataEvent, DataResponse,
37        data::{
38            BarsResponse, BookResponse, FundingRatesResponse, InstrumentResponse,
39            InstrumentsResponse, RequestBars, RequestBookSnapshot, RequestFundingRates,
40            RequestInstrument, RequestInstruments, RequestTrades, SubscribeBars,
41            SubscribeBookDeltas, SubscribeFundingRates, SubscribeIndexPrices, SubscribeInstrument,
42            SubscribeInstrumentClose, SubscribeInstrumentStatus, SubscribeInstruments,
43            SubscribeMarkPrices, SubscribeQuotes, SubscribeTrades, TradesResponse, UnsubscribeBars,
44            UnsubscribeBookDeltas, UnsubscribeFundingRates, UnsubscribeIndexPrices,
45            UnsubscribeInstrument, UnsubscribeInstrumentClose, UnsubscribeInstrumentStatus,
46            UnsubscribeInstruments, UnsubscribeMarkPrices, UnsubscribeQuotes, UnsubscribeTrades,
47        },
48    },
49};
50use nautilus_core::{
51    AtomicMap,
52    datetime::datetime_to_unix_nanos,
53    nanos::UnixNanos,
54    time::{AtomicTime, get_atomic_clock_realtime},
55};
56use nautilus_live::{
57    SocketControl,
58    task::{TaskGroup, TaskGroupGuard},
59};
60use nautilus_model::{
61    data::{Data, FundingRateUpdate, InstrumentStatus, MarkPriceUpdate},
62    enums::{BookType, MarketStatusAction},
63    identifiers::{ClientId, InstrumentId, Venue},
64    instruments::{Instrument, InstrumentAny},
65    types::Price,
66};
67use parking_lot::Mutex;
68use tokio_util::sync::CancellationToken;
69use ustr::Ustr;
70
71use crate::{
72    common::{
73        auth::run_auth_token_refresh,
74        consts::{AX_AUTH_TOKEN_TTL_SECS, AX_FUNDING_RATE_LOOKBACK_DAYS, AX_VENUE},
75        credential::Credential,
76        enums::{AxCandleWidth, AxInstrumentState, AxMarketDataLevel},
77        parse::{ax_timestamp_stn_to_unix_nanos, map_bar_spec_to_candle_width},
78    },
79    config::AxDataClientConfig,
80    http::client::AxHttpClient,
81    websocket::{
82        data::{
83            client::{AxMdWebSocketClient, AxWsClientError, SymbolDataTypes},
84            parse::{
85                parse_book_l1_quote, parse_book_l2_deltas, parse_book_l2_quote,
86                parse_book_l3_deltas, parse_book_l3_quote, parse_candle_bar, parse_trade_tick,
87            },
88        },
89        messages::{AxDataWsMessage, AxMdCandle, AxMdMessage},
90    },
91};
92
93/// AX Exchange data client for live market data streaming and historical data requests.
94///
95/// This client integrates with the Nautilus DataEngine to provide:
96/// - Real-time market data via WebSocket subscriptions
97/// - Historical data via REST API requests
98/// - Automatic instrument discovery and caching
99/// - Connection lifecycle management
100#[derive(Debug)]
101pub struct AxDataClient {
102    client_id: ClientId,
103    config: AxDataClientConfig,
104    http_client: AxHttpClient,
105    ws_client: AxMdWebSocketClient,
106    is_connected: Arc<AtomicBool>,
107    cancellation_token: CancellationToken,
108    session_tasks: TaskGroup,
109    pending_tasks: TaskGroup,
110    shutdown_errors: Vec<String>,
111    data_sender: EventSender<DataEvent>,
112    instruments: Arc<AtomicMap<Ustr, InstrumentAny>>,
113    clock: &'static AtomicTime,
114    funding_rate_cancellations: AHashMap<InstrumentId, CancellationToken>,
115    funding_rate_cache: Arc<Mutex<AHashMap<InstrumentId, FundingRateUpdate>>>,
116}
117
118impl AxDataClient {
119    /// Creates a new [`AxDataClient`] instance.
120    ///
121    /// # Errors
122    ///
123    /// Returns an error if the data event sender cannot be obtained.
124    pub fn new(
125        client_id: ClientId,
126        config: AxDataClientConfig,
127        http_client: AxHttpClient,
128        ws_client: AxMdWebSocketClient,
129    ) -> anyhow::Result<Self> {
130        let clock = get_atomic_clock_realtime();
131        let data_sender = get_data_event_sender();
132        let ws_client = ws_client.with_socket_control(SocketControl::new(
133            client_id,
134            Some(*AX_VENUE),
135            "architect-ax-data-streams",
136        ));
137
138        // Share instruments cache with HTTP client
139        let instruments = http_client.instruments_cache.clone();
140
141        let session_tasks = TaskGroup::new();
142        let pending_tasks = TaskGroup::new();
143
144        Ok(Self {
145            client_id,
146            config,
147            http_client,
148            ws_client,
149            is_connected: Arc::new(AtomicBool::new(false)),
150            cancellation_token: CancellationToken::new(),
151            session_tasks,
152            pending_tasks,
153            shutdown_errors: Vec::new(),
154            data_sender,
155            instruments,
156            clock,
157            funding_rate_cancellations: AHashMap::new(),
158            funding_rate_cache: Arc::new(Mutex::new(AHashMap::new())),
159        })
160    }
161
162    /// Returns the venue for this data client.
163    #[must_use]
164    pub fn venue(&self) -> Venue {
165        *AX_VENUE
166    }
167
168    fn map_book_type_to_market_data_level(book_type: BookType) -> AxMarketDataLevel {
169        match book_type {
170            BookType::L3_MBO => AxMarketDataLevel::Level3,
171            BookType::L1_MBP | BookType::L2_MBP => AxMarketDataLevel::Level2,
172        }
173    }
174
175    /// Returns a reference to the instruments cache.
176    #[must_use]
177    pub fn instruments(&self) -> &Arc<AtomicMap<Ustr, InstrumentAny>> {
178        &self.instruments
179    }
180
181    /// Spawns a message handler task to forward WebSocket data to the DataEngine.
182    fn spawn_message_handler(&mut self) -> anyhow::Result<()> {
183        let stream = self.ws_client.stream();
184        let data_sender = self.data_sender.clone();
185        let cancellation_token = self.cancellation_token.clone();
186        let is_connected = Arc::clone(&self.is_connected);
187        let instruments = Arc::clone(&self.instruments);
188        let symbol_data_types = self.ws_client.symbol_data_types();
189        let status_invalidations = self.ws_client.status_invalidations();
190        let clock = self.clock;
191
192        self.session_tasks.spawn(async move {
193            tokio::pin!(stream);
194
195            let mut book_sequences: AHashMap<Ustr, u64> = AHashMap::new();
196            let mut candle_cache: AHashMap<(Ustr, AxCandleWidth), AxMdCandle> = AHashMap::new();
197            let mut instrument_states: AHashMap<Ustr, AxInstrumentState> = AHashMap::new();
198
199            loop {
200                tokio::select! {
201                    () = cancellation_token.cancelled() => {
202                        log::debug!("Message handler cancelled");
203                        break;
204                    }
205                    msg = stream.next() => {
206                        match msg {
207                            Some(ws_msg) => {
208                                drain_status_invalidations(
209                                    &status_invalidations,
210                                    &mut instrument_states,
211                                );
212
213                                handle_ws_message(
214                                    ws_msg,
215                                    &data_sender,
216                                    &instruments,
217                                    &symbol_data_types,
218                                    &mut book_sequences,
219                                    &mut candle_cache,
220                                    &mut instrument_states,
221                                    clock,
222                                );
223                            }
224                            None => {
225                                log::debug!("WebSocket stream ended");
226                                is_connected.store(false, Ordering::Release);
227                                break;
228                            }
229                        }
230                    }
231                }
232            }
233        })?;
234        Ok(())
235    }
236
237    fn spawn_instrument_refresh(&self) -> anyhow::Result<()> {
238        let minutes = self.config.update_instruments_interval_mins;
239        if minutes == 0 {
240            return Ok(());
241        }
242
243        let interval = Duration::from_secs(minutes.saturating_mul(60));
244        let cancellation = self.cancellation_token.clone();
245        let instruments_cache = Arc::clone(&self.instruments);
246        let http_client = self.http_client.clone();
247        let data_sender = self.data_sender.clone();
248        let client_id = self.client_id;
249
250        self.session_tasks.spawn(async move {
251            loop {
252                let sleep = tokio::time::sleep(interval);
253                tokio::pin!(sleep);
254                tokio::select! {
255                    () = cancellation.cancelled() => {
256                        log::debug!("Instrument refresh task cancelled");
257                        break;
258                    }
259                    () = &mut sleep => {
260                        match http_client.request_instruments(None, None).await {
261                            Ok(instruments) => {
262                                for inst in &instruments {
263                                    instruments_cache.insert(inst.symbol().inner(), inst.clone());
264
265                                    if let Err(e) = data_sender
266                                        .send(DataEvent::Instrument(inst.clone()))
267                                    {
268                                        log::warn!("Failed to send refreshed instrument: {e}");
269                                    }
270                                }
271                                http_client.cache_instruments(&instruments);
272                                log::debug!(
273                                    "Instruments refreshed: client_id={client_id}, count={}",
274                                    instruments.len(),
275                                );
276                            }
277                            Err(e) => {
278                                log::warn!("Failed to refresh instruments: client_id={client_id}, error={e:?}");
279                            }
280                        }
281                    }
282                }
283            }
284        })?;
285        Ok(())
286    }
287
288    #[expect(
289        clippy::unnecessary_wraps,
290        reason = "callers forward Result to trait methods"
291    )]
292    fn ws_symbol_op<F, Fut>(
293        &self,
294        instrument_id: InstrumentId,
295        op: F,
296        context: &'static str,
297    ) -> anyhow::Result<()>
298    where
299        F: FnOnce(AxMdWebSocketClient, String) -> Fut + Send + 'static,
300        Fut: Future<Output = Result<(), AxWsClientError>> + Send,
301    {
302        let symbol = instrument_id.symbol.to_string();
303        log::debug!("{context} for {symbol}");
304
305        let ws = self.ws_client.clone();
306        self.spawn_ws(
307            async move { op(ws, symbol).await.map_err(|e| anyhow::anyhow!(e)) },
308            context,
309        );
310
311        Ok(())
312    }
313
314    fn spawn_ws<F>(&self, fut: F, context: &'static str)
315    where
316        F: Future<Output = anyhow::Result<()>> + Send + 'static,
317    {
318        let future = async move {
319            if let Err(e) = fut.await {
320                log::error!("{context}: {e:?}");
321            }
322        };
323
324        if let Err(e) = self.pending_tasks.spawn(future) {
325            log::warn!("Skipping AX {context} after shutdown began: {e}");
326        }
327    }
328
329    fn spawn_task<F>(&self, fut: F)
330    where
331        F: Future<Output = ()> + Send + 'static,
332    {
333        if let Err(e) = self.pending_tasks.spawn(fut) {
334            log::warn!("Skipping AX data task after shutdown began: {e}");
335        }
336    }
337
338    fn abort_pending_tasks(&self) {
339        self.pending_tasks.begin_shutdown();
340    }
341
342    fn abort_all_tasks(&self) {
343        self.cancellation_token.cancel();
344        self.session_tasks.begin_shutdown();
345        self.abort_pending_tasks();
346        self.ws_client.begin_shutdown();
347
348        for cancellation in self.funding_rate_cancellations.values() {
349            cancellation.cancel();
350        }
351    }
352
353    async fn finish_all_tasks(&mut self) -> anyhow::Result<()> {
354        self.pending_tasks.begin_shutdown();
355        self.session_tasks.begin_shutdown();
356        let (pending_result, session_result) = tokio::join!(
357            self.pending_tasks
358                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
359            self.session_tasks
360                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
361        );
362        self.funding_rate_cancellations.clear();
363
364        pending_result.map_err(|e| anyhow::anyhow!("Failed to terminate AX data tasks: {e}"))?;
365        session_result
366            .map_err(|e| anyhow::anyhow!("Failed to terminate AX data session tasks: {e}"))?;
367        Ok(())
368    }
369
370    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
371        self.abort_all_tasks();
372
373        if let Err(e) = self.ws_client.close().await {
374            self.shutdown_errors.push(e.to_string());
375        }
376
377        if let Err(e) = self.finish_all_tasks().await {
378            self.shutdown_errors.push(e.to_string());
379        }
380        self.is_connected.store(false, Ordering::Release);
381
382        if !self.shutdown_errors.is_empty() {
383            anyhow::bail!(std::mem::take(&mut self.shutdown_errors).join("; "));
384        }
385        Ok(())
386    }
387}
388
389#[async_trait(?Send)]
390impl DataClient for AxDataClient {
391    fn client_id(&self) -> ClientId {
392        self.client_id
393    }
394
395    fn venue(&self) -> Option<Venue> {
396        Some(*AX_VENUE)
397    }
398
399    fn start(&mut self) -> anyhow::Result<()> {
400        log::debug!("Starting {}", self.client_id);
401        Ok(())
402    }
403
404    fn stop(&mut self) -> anyhow::Result<()> {
405        log::debug!("Stopping {}", self.client_id);
406
407        self.abort_all_tasks();
408        self.is_connected.store(false, Ordering::Release);
409        Ok(())
410    }
411
412    fn reset(&mut self) -> anyhow::Result<()> {
413        log::debug!("Resetting {}", self.client_id);
414
415        self.abort_all_tasks();
416        self.is_connected.store(false, Ordering::Release);
417        self.funding_rate_cache.lock().clear();
418        Ok(())
419    }
420
421    fn dispose(&mut self) -> anyhow::Result<()> {
422        log::debug!("Disposing {}", self.client_id);
423
424        self.abort_all_tasks();
425        self.is_connected.store(false, Ordering::Release);
426        Ok(())
427    }
428
429    fn is_connected(&self) -> bool {
430        self.is_connected.load(Ordering::Acquire)
431    }
432
433    fn is_disconnected(&self) -> bool {
434        !self.is_connected()
435    }
436
437    async fn connect(&mut self) -> anyhow::Result<()> {
438        if self.is_connected()
439            && !self.cancellation_token.is_cancelled()
440            && self.pending_tasks.is_open()
441            && self.session_tasks.is_open()
442        {
443            log::debug!("Already connected {}", self.client_id);
444            return Ok(());
445        }
446
447        log::info!("Connecting {}", self.client_id);
448
449        if self.cancellation_token.is_cancelled()
450            || !self.pending_tasks.is_open()
451            || !self.session_tasks.is_open()
452            || !self.funding_rate_cancellations.is_empty()
453        {
454            self.teardown_partial_connect().await?;
455            self.session_tasks
456                .start_generation()
457                .map_err(|e| anyhow::anyhow!("Failed to start AX data session generation: {e}"))?;
458            self.pending_tasks
459                .start_generation()
460                .map_err(|e| anyhow::anyhow!("Failed to start AX data task generation: {e}"))?;
461            self.cancellation_token = CancellationToken::new();
462        }
463        let cancellation_token = self.cancellation_token.clone();
464        let ws_client = self.ws_client.clone();
465        let setup_guard =
466            TaskGroupGuard::new(&[&self.session_tasks, &self.pending_tasks], move || {
467                cancellation_token.cancel();
468                ws_client.begin_shutdown();
469            });
470
471        let credential = if self.config.has_api_credentials() {
472            let credential = Credential::resolve(
473                self.config.api_key.clone().map(|value| value.into_inner()),
474                self.config
475                    .api_secret
476                    .clone()
477                    .map(|value| value.into_inner()),
478            )
479            .context("API credentials not configured")?;
480
481            let token = self
482                .http_client
483                .authenticate(
484                    credential.api_key(),
485                    credential.api_secret(),
486                    AX_AUTH_TOKEN_TTL_SECS,
487                )
488                .await
489                .context("Failed to authenticate with Ax")?;
490            log::debug!("Authenticated with Ax");
491            self.ws_client.set_auth_token(token);
492
493            // Only an authenticated client can read fee rates, and a data client may
494            // legitimately run without credentials.
495            self.http_client
496                .request_account_fees()
497                .await
498                .context("Failed to resolve Ax account fee rates")?;
499
500            Some(credential)
501        } else {
502            log::debug!("No Ax credentials configured, instruments will report zero fees");
503            None
504        };
505
506        let instruments = self
507            .http_client
508            .request_instruments(None, None)
509            .await
510            .context("Failed to fetch instruments")?;
511
512        for instrument in &instruments {
513            self.instruments
514                .insert(instrument.symbol().inner(), instrument.clone());
515
516            if let Err(e) = self
517                .data_sender
518                .send(DataEvent::Instrument(instrument.clone()))
519            {
520                log::warn!("Failed to send instrument: {e}");
521            }
522        }
523        self.http_client.cache_instruments(&instruments);
524        log::debug!(
525            "Cached {} instruments",
526            self.http_client.get_cached_symbols().len()
527        );
528
529        self.ws_client
530            .connect()
531            .await
532            .context("Failed to connect WebSocket")?;
533        log::debug!("WebSocket connected");
534
535        let session_result = async {
536            self.spawn_message_handler()?;
537            self.spawn_instrument_refresh()?;
538
539            if let Some(credential) = credential {
540                let ws_client = self.ws_client.clone();
541                self.session_tasks.spawn(run_auth_token_refresh(
542                    self.http_client.clone(),
543                    credential,
544                    move |token| ws_client.update_auth_token(token),
545                ))?;
546            }
547            Ok::<(), anyhow::Error>(())
548        }
549        .await;
550
551        if let Err(e) = session_result {
552            if let Err(teardown_error) = self.teardown_partial_connect().await {
553                return Err(e.context(format!("AX data startup teardown failed: {teardown_error}")));
554            }
555            return Err(e);
556        }
557
558        self.is_connected.store(true, Ordering::Release);
559        setup_guard.disarm();
560        log::info!("Connected {}", self.client_id);
561
562        Ok(())
563    }
564
565    async fn disconnect(&mut self) -> anyhow::Result<()> {
566        log::info!("Disconnecting {}", self.client_id);
567
568        self.abort_all_tasks();
569        let ws_result = self.ws_client.close().await;
570        let tasks_result = self.finish_all_tasks().await;
571        self.funding_rate_cache.lock().clear();
572
573        self.is_connected.store(false, Ordering::Release);
574        log::info!("Disconnected {}", self.client_id);
575
576        ws_result?;
577        tasks_result
578    }
579
580    fn subscribe_instruments(&mut self, _cmd: SubscribeInstruments) -> anyhow::Result<()> {
581        // AX does not have a real-time instruments channel; instruments are fetched via HTTP
582        log::debug!("Instruments subscription not applicable for AX (use request_instruments)");
583        Ok(())
584    }
585
586    fn subscribe_instrument(&mut self, _cmd: SubscribeInstrument) -> anyhow::Result<()> {
587        // AX does not have a real-time instrument channel; instruments are fetched via HTTP
588        log::debug!("Instrument subscription not applicable for AX (use request_instrument)");
589        Ok(())
590    }
591
592    fn subscribe_book_deltas(&mut self, cmd: SubscribeBookDeltas) -> anyhow::Result<()> {
593        let symbol = cmd.instrument_id.symbol.to_string();
594        let level = Self::map_book_type_to_market_data_level(cmd.book_type);
595        if cmd.book_type == BookType::L1_MBP {
596            log::warn!(
597                "Book type L1_MBP not supported by AX for deltas, downgrading {symbol} to LEVEL_2"
598            );
599        }
600        log::debug!("Subscribing to book deltas for {symbol} at {level:?}");
601
602        let ws = self.ws_client.clone();
603        self.spawn_ws(
604            async move {
605                ws.subscribe_book_deltas(&symbol, level)
606                    .await
607                    .map_err(|e| anyhow::anyhow!(e))
608            },
609            "subscribe book deltas",
610        );
611
612        Ok(())
613    }
614
615    fn subscribe_quotes(&mut self, cmd: SubscribeQuotes) -> anyhow::Result<()> {
616        self.ws_symbol_op(
617            cmd.instrument_id,
618            |ws, s| async move { ws.subscribe_quotes(&s).await },
619            "Subscribing to quotes",
620        )
621    }
622
623    fn subscribe_trades(&mut self, cmd: SubscribeTrades) -> anyhow::Result<()> {
624        self.ws_symbol_op(
625            cmd.instrument_id,
626            |ws, s| async move { ws.subscribe_trades(&s).await },
627            "Subscribing to trades",
628        )
629    }
630
631    fn subscribe_mark_prices(&mut self, cmd: SubscribeMarkPrices) -> anyhow::Result<()> {
632        self.ws_symbol_op(
633            cmd.instrument_id,
634            |ws, s| async move { ws.subscribe_mark_prices(&s).await },
635            "Subscribing to mark prices",
636        )
637    }
638
639    fn subscribe_index_prices(&mut self, _cmd: SubscribeIndexPrices) -> anyhow::Result<()> {
640        log::warn!("Index prices not supported by AX Exchange");
641        Ok(())
642    }
643
644    fn subscribe_bars(&mut self, cmd: SubscribeBars) -> anyhow::Result<()> {
645        let bar_type = cmd.bar_type;
646        let symbol = bar_type.instrument_id().symbol.to_string();
647        let width = map_bar_spec_to_candle_width(&bar_type.spec())?;
648        log::debug!("Subscribing to bars for {bar_type} (width: {width:?})");
649
650        let ws = self.ws_client.clone();
651        self.spawn_ws(
652            async move {
653                ws.subscribe_candles(&symbol, width)
654                    .await
655                    .map_err(|e| anyhow::anyhow!(e))
656            },
657            "subscribe bars",
658        );
659
660        Ok(())
661    }
662
663    fn subscribe_funding_rates(&mut self, cmd: SubscribeFundingRates) -> anyhow::Result<()> {
664        let poll_interval_mins = self.config.funding_rate_poll_interval_mins.max(1);
665
666        // Use 7-day lookback to capture latest rate across weekends/holidays
667        let lookback = SignedDuration::from_hours(24 * (AX_FUNDING_RATE_LOOKBACK_DAYS));
668
669        let instrument_id = cmd.instrument_id;
670
671        if self.funding_rate_cancellations.contains_key(&instrument_id) {
672            log::debug!("Already subscribed to funding rates for {instrument_id}");
673            return Ok(());
674        }
675
676        log::debug!("Subscribing to funding rates for {instrument_id} (HTTP polling)");
677
678        let http = self.http_client.clone();
679        let sender = self.data_sender.clone();
680        let symbol = instrument_id.symbol.inner();
681        let cancellation = self.cancellation_token.child_token();
682        let task_cancellation = cancellation.clone();
683        let cache = Arc::clone(&self.funding_rate_cache);
684        let clock = self.clock;
685
686        self.session_tasks.spawn(async move {
687            // First tick fires immediately for initial emission
688            let mut interval = tokio::time::interval(Duration::from_mins(poll_interval_mins));
689
690            loop {
691                tokio::select! {
692                    () = task_cancellation.cancelled() => {
693                        log::debug!("Funding rate polling cancelled for {symbol}");
694                        break;
695                    }
696                    _ = interval.tick() => {
697                        let now: Timestamp = clock.get_time_ns().into();
698                        let start = now - lookback;
699
700                        match http.request_funding_rates(instrument_id, Some(start), Some(now)).await {
701                            Ok(funding_rates) => {
702                                if funding_rates.is_empty() {
703                                    log::warn!(
704                                        "No funding rates returned for {symbol}"
705                                    );
706                                } else if let Some(update) = funding_rates.last() {
707                                    // Only emit if rate changed
708                                    let should_emit = cache.lock()
709                                        .get(&instrument_id) != Some(update);
710
711                                    if should_emit {
712                                        log::debug!(
713                                            "Funding rate for {symbol}: {}",
714                                            update.rate,
715                                        );
716                                        let update = *update;
717                                        cache.lock()
718                                            .insert(instrument_id, update);
719
720                                        if let Err(e) = sender.send(
721                                            DataEvent::FundingRate(update),
722                                        ) {
723                                            log::error!(
724                                                "Failed to send funding rate for {symbol}: {e}"
725                                            );
726                                        }
727                                    }
728                                }
729                            }
730                            Err(e) => {
731                                log::error!(
732                                    "Failed to poll funding rates for {symbol}: {e}"
733                                );
734                            }
735                        }
736                    }
737                }
738            }
739        })?;
740
741        self.funding_rate_cancellations
742            .insert(instrument_id, cancellation);
743        Ok(())
744    }
745
746    fn subscribe_instrument_status(
747        &mut self,
748        cmd: SubscribeInstrumentStatus,
749    ) -> anyhow::Result<()> {
750        self.ws_symbol_op(
751            cmd.instrument_id,
752            |ws, s| async move { ws.subscribe_instrument_status(&s).await },
753            "Subscribing to instrument status",
754        )
755    }
756
757    fn subscribe_instrument_close(&mut self, _cmd: SubscribeInstrumentClose) -> anyhow::Result<()> {
758        log::warn!("Instrument close not supported by AX Exchange");
759        Ok(())
760    }
761
762    fn unsubscribe_instruments(&mut self, _cmd: &UnsubscribeInstruments) -> anyhow::Result<()> {
763        Ok(())
764    }
765
766    fn unsubscribe_instrument(&mut self, _cmd: &UnsubscribeInstrument) -> anyhow::Result<()> {
767        Ok(())
768    }
769
770    fn unsubscribe_book_deltas(&mut self, cmd: &UnsubscribeBookDeltas) -> anyhow::Result<()> {
771        self.ws_symbol_op(
772            cmd.instrument_id,
773            |ws, s| async move { ws.unsubscribe_book_deltas(&s).await },
774            "Unsubscribing from book deltas",
775        )
776    }
777
778    fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
779        self.ws_symbol_op(
780            cmd.instrument_id,
781            |ws, s| async move { ws.unsubscribe_quotes(&s).await },
782            "Unsubscribing from quotes",
783        )
784    }
785
786    fn unsubscribe_trades(&mut self, cmd: &UnsubscribeTrades) -> anyhow::Result<()> {
787        self.ws_symbol_op(
788            cmd.instrument_id,
789            |ws, s| async move { ws.unsubscribe_trades(&s).await },
790            "Unsubscribing from trades",
791        )
792    }
793
794    fn unsubscribe_mark_prices(&mut self, cmd: &UnsubscribeMarkPrices) -> anyhow::Result<()> {
795        self.ws_symbol_op(
796            cmd.instrument_id,
797            |ws, s| async move { ws.unsubscribe_mark_prices(&s).await },
798            "Unsubscribing from mark prices",
799        )
800    }
801
802    fn unsubscribe_index_prices(&mut self, _cmd: &UnsubscribeIndexPrices) -> anyhow::Result<()> {
803        Ok(())
804    }
805
806    fn unsubscribe_bars(&mut self, cmd: &UnsubscribeBars) -> anyhow::Result<()> {
807        let bar_type = cmd.bar_type;
808        let symbol = bar_type.instrument_id().symbol.to_string();
809        let width = map_bar_spec_to_candle_width(&bar_type.spec())?;
810        log::debug!("Unsubscribing from bars for {bar_type}");
811
812        let ws = self.ws_client.clone();
813        self.spawn_ws(
814            async move {
815                ws.unsubscribe_candles(&symbol, width)
816                    .await
817                    .map_err(|e| anyhow::anyhow!(e))
818            },
819            "unsubscribe bars",
820        );
821
822        Ok(())
823    }
824
825    fn unsubscribe_funding_rates(&mut self, cmd: &UnsubscribeFundingRates) -> anyhow::Result<()> {
826        let instrument_id = cmd.instrument_id;
827
828        if let Some(cancellation) = self.funding_rate_cancellations.remove(&instrument_id) {
829            log::debug!("Unsubscribing from funding rates for {instrument_id}");
830            cancellation.cancel();
831            self.funding_rate_cache.lock().remove(&instrument_id);
832        } else {
833            log::debug!("Not subscribed to funding rates for {instrument_id}");
834        }
835
836        Ok(())
837    }
838
839    fn unsubscribe_instrument_status(
840        &mut self,
841        cmd: &UnsubscribeInstrumentStatus,
842    ) -> anyhow::Result<()> {
843        self.ws_symbol_op(
844            cmd.instrument_id,
845            |ws, s| async move { ws.unsubscribe_instrument_status(&s).await },
846            "Unsubscribing from instrument status",
847        )
848    }
849
850    fn unsubscribe_instrument_close(
851        &mut self,
852        _cmd: &UnsubscribeInstrumentClose,
853    ) -> anyhow::Result<()> {
854        Ok(())
855    }
856
857    fn request_instruments(&self, request: RequestInstruments) -> anyhow::Result<()> {
858        let http = self.http_client.clone();
859        let instruments_cache = Arc::clone(&self.instruments);
860        let sender = self.data_sender.clone();
861        let cancel = self.cancellation_token.clone();
862        let request_id = request.request_id;
863        let client_id = request.client_id.unwrap_or(self.client_id);
864        let venue = *AX_VENUE;
865        let start_nanos = datetime_to_unix_nanos(request.start);
866        let end_nanos = datetime_to_unix_nanos(request.end);
867        let params = request.params;
868        let clock = self.clock;
869
870        self.spawn_task(async move {
871            match http.request_instruments(None, None).await {
872                Ok(instruments) => {
873                    if cancel.is_cancelled() {
874                        return;
875                    }
876                    log::debug!("Fetched {} instruments from Ax", instruments.len());
877                    for inst in &instruments {
878                        instruments_cache.insert(inst.symbol().inner(), inst.clone());
879                    }
880                    http.cache_instruments(&instruments);
881
882                    let response = DataResponse::Instruments(InstrumentsResponse::new(
883                        request_id,
884                        client_id,
885                        venue,
886                        instruments,
887                        start_nanos,
888                        end_nanos,
889                        clock.get_time_ns(),
890                        params,
891                    ));
892
893                    if let Err(e) = sender.send(DataEvent::Response(response)) {
894                        log::error!("Failed to send instruments response: {e}");
895                    }
896                }
897                Err(e) => {
898                    log::error!("Failed to request instruments: {e}");
899                }
900            }
901        });
902
903        Ok(())
904    }
905
906    fn request_instrument(&self, request: RequestInstrument) -> anyhow::Result<()> {
907        let http = self.http_client.clone();
908        let instruments_cache = Arc::clone(&self.instruments);
909        let sender = self.data_sender.clone();
910        let cancel = self.cancellation_token.clone();
911        let request_id = request.request_id;
912        let client_id = request.client_id.unwrap_or(self.client_id);
913        let instrument_id = request.instrument_id;
914        let symbol = instrument_id.symbol.inner();
915        let start_nanos = datetime_to_unix_nanos(request.start);
916        let end_nanos = datetime_to_unix_nanos(request.end);
917        let params = request.params;
918        let clock = self.clock;
919
920        self.spawn_task(async move {
921            match http.request_instrument(symbol, None, None).await {
922                Ok(instrument) => {
923                    if cancel.is_cancelled() {
924                        return;
925                    }
926                    log::debug!("Fetched instrument {symbol} from Ax");
927                    instruments_cache.insert(symbol, instrument.clone());
928                    http.cache_instrument(instrument.clone());
929
930                    let response = DataResponse::Instrument(Box::new(InstrumentResponse::new(
931                        request_id,
932                        client_id,
933                        instrument_id,
934                        instrument,
935                        start_nanos,
936                        end_nanos,
937                        clock.get_time_ns(),
938                        params,
939                    )));
940
941                    if let Err(e) = sender.send(DataEvent::Response(response)) {
942                        log::error!("Failed to send instrument response: {e}");
943                    }
944                }
945                Err(e) => {
946                    log::error!("Failed to request instrument {symbol}: {e}");
947                }
948            }
949        });
950
951        Ok(())
952    }
953
954    fn request_book_snapshot(&self, request: RequestBookSnapshot) -> anyhow::Result<()> {
955        let http = self.http_client.clone();
956        let sender = self.data_sender.clone();
957        let cancel = self.cancellation_token.clone();
958        let request_id = request.request_id;
959        let client_id = request.client_id.unwrap_or(self.client_id);
960        let instrument_id = request.instrument_id;
961        let symbol = instrument_id.symbol.inner();
962        let depth = request.depth.map(|n| n.get());
963        let params = request.params;
964        let clock = self.clock;
965
966        self.spawn_task(async move {
967            match http.request_book_snapshot(symbol, depth).await {
968                Ok(book) => {
969                    if cancel.is_cancelled() {
970                        return;
971                    }
972                    log::debug!(
973                        "Fetched book snapshot for {symbol} ({} bids, {} asks)",
974                        book.bids(None).count(),
975                        book.asks(None).count(),
976                    );
977
978                    let response = DataResponse::Book(BookResponse::new(
979                        request_id,
980                        client_id,
981                        instrument_id,
982                        book,
983                        None,
984                        None,
985                        clock.get_time_ns(),
986                        params,
987                    ));
988
989                    if let Err(e) = sender.send(DataEvent::Response(response)) {
990                        log::error!("Failed to send book snapshot response: {e}");
991                    }
992                }
993                Err(e) => {
994                    log::error!("Failed to request book snapshot for {symbol}: {e}");
995                }
996            }
997        });
998
999        Ok(())
1000    }
1001
1002    fn request_trades(&self, request: RequestTrades) -> anyhow::Result<()> {
1003        let http = self.http_client.clone();
1004        let sender = self.data_sender.clone();
1005        let cancel = self.cancellation_token.clone();
1006        let request_id = request.request_id;
1007        let client_id = request.client_id.unwrap_or(self.client_id);
1008        let instrument_id = request.instrument_id;
1009        let symbol = instrument_id.symbol.inner();
1010        let limit = request.limit.map(|n| n.get() as i32);
1011        let start_nanos = datetime_to_unix_nanos(request.start);
1012        let end_nanos = datetime_to_unix_nanos(request.end);
1013        let params = request.params;
1014        let clock = self.clock;
1015
1016        self.spawn_task(async move {
1017            match http
1018                .request_trade_ticks(symbol, limit, start_nanos, end_nanos)
1019                .await
1020            {
1021                Ok(ticks) => {
1022                    if cancel.is_cancelled() {
1023                        return;
1024                    }
1025                    log::debug!("Fetched {} trades for {symbol}", ticks.len());
1026
1027                    let response = DataResponse::Trades(TradesResponse::new(
1028                        request_id,
1029                        client_id,
1030                        instrument_id,
1031                        ticks,
1032                        start_nanos,
1033                        end_nanos,
1034                        clock.get_time_ns(),
1035                        params,
1036                    ));
1037
1038                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1039                        log::error!("Failed to send trades response: {e}");
1040                    }
1041                }
1042                Err(e) => {
1043                    log::error!("Failed to request trades for {symbol}: {e}");
1044                }
1045            }
1046        });
1047
1048        Ok(())
1049    }
1050
1051    fn request_bars(&self, request: RequestBars) -> anyhow::Result<()> {
1052        let http = self.http_client.clone();
1053        let sender = self.data_sender.clone();
1054        let request_id = request.request_id;
1055        let client_id = request.client_id.unwrap_or(self.client_id);
1056        let bar_type = request.bar_type;
1057        let symbol = bar_type.instrument_id().symbol.inner();
1058        let start = request.start;
1059        let end = request.end;
1060        let start_nanos = datetime_to_unix_nanos(start);
1061        let end_nanos = datetime_to_unix_nanos(end);
1062        let params = request.params;
1063        let clock = self.clock;
1064        let width = match map_bar_spec_to_candle_width(&bar_type.spec()) {
1065            Ok(w) => w,
1066            Err(e) => {
1067                log::error!("Failed to map bar type {bar_type}: {e}");
1068                return Err(e);
1069            }
1070        };
1071
1072        let cancel = self.cancellation_token.clone();
1073
1074        self.spawn_task(async move {
1075            match http.request_bars(symbol, start, end, width).await {
1076                Ok(bars) => {
1077                    if cancel.is_cancelled() {
1078                        return;
1079                    }
1080                    log::debug!("Fetched {} bars for {symbol}", bars.len());
1081
1082                    let response = DataResponse::Bars(BarsResponse::new(
1083                        request_id,
1084                        client_id,
1085                        bar_type,
1086                        bars,
1087                        start_nanos,
1088                        end_nanos,
1089                        clock.get_time_ns(),
1090                        params,
1091                    ));
1092
1093                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1094                        log::error!("Failed to send bars response: {e}");
1095                    }
1096                }
1097                Err(e) => {
1098                    log::error!("Failed to request bars for {symbol}: {e}");
1099                }
1100            }
1101        });
1102
1103        Ok(())
1104    }
1105
1106    fn request_funding_rates(&self, request: RequestFundingRates) -> anyhow::Result<()> {
1107        let http = self.http_client.clone();
1108        let sender = self.data_sender.clone();
1109        let cancel = self.cancellation_token.clone();
1110        let request_id = request.request_id;
1111        let client_id = request.client_id.unwrap_or(self.client_id);
1112        let instrument_id = request.instrument_id;
1113        let symbol = instrument_id.symbol.inner();
1114        let start = request.start;
1115        let end = request.end;
1116        let start_nanos = datetime_to_unix_nanos(start);
1117        let end_nanos = datetime_to_unix_nanos(end);
1118        let params = request.params;
1119        let clock = self.clock;
1120
1121        self.spawn_task(async move {
1122            match http.request_funding_rates(instrument_id, start, end).await {
1123                Ok(funding_rates) => {
1124                    if cancel.is_cancelled() {
1125                        return;
1126                    }
1127                    log::debug!("Fetched {} funding rates for {symbol}", funding_rates.len());
1128
1129                    let ts_init = clock.get_time_ns();
1130                    let response = DataResponse::FundingRates(FundingRatesResponse::new(
1131                        request_id,
1132                        client_id,
1133                        instrument_id,
1134                        funding_rates,
1135                        start_nanos,
1136                        end_nanos,
1137                        ts_init,
1138                        params,
1139                    ));
1140
1141                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1142                        log::error!("Failed to send funding rates response: {e}");
1143                    }
1144                }
1145                Err(e) => {
1146                    log::error!("Failed to request funding rates for {symbol}: {e}");
1147                }
1148            }
1149        });
1150
1151        Ok(())
1152    }
1153}
1154
1155fn drain_status_invalidations(
1156    invalidations: &Arc<Mutex<AHashSet<Ustr>>>,
1157    instrument_states: &mut AHashMap<Ustr, AxInstrumentState>,
1158) {
1159    for symbol in invalidations.lock().drain() {
1160        instrument_states.remove(&symbol);
1161    }
1162}
1163
1164#[expect(clippy::too_many_arguments)]
1165fn handle_ws_message(
1166    msg: AxDataWsMessage,
1167    sender: &EventSender<DataEvent>,
1168    instruments: &Arc<AtomicMap<Ustr, InstrumentAny>>,
1169    symbol_data_types: &Arc<AtomicMap<String, SymbolDataTypes>>,
1170    book_sequences: &mut AHashMap<Ustr, u64>,
1171    candle_cache: &mut AHashMap<(Ustr, AxCandleWidth), AxMdCandle>,
1172    instrument_states: &mut AHashMap<Ustr, AxInstrumentState>,
1173    clock: &'static AtomicTime,
1174) {
1175    match msg {
1176        AxDataWsMessage::Reconnected => {
1177            candle_cache.clear();
1178            instrument_states.clear();
1179            log::info!("WebSocket reconnected");
1180        }
1181        AxDataWsMessage::CandleUnsubscribed { symbol, width } => {
1182            candle_cache.remove(&(symbol, width));
1183        }
1184        AxDataWsMessage::MdMessage(md_msg) => {
1185            handle_md_message(
1186                md_msg,
1187                sender,
1188                instruments,
1189                symbol_data_types,
1190                book_sequences,
1191                candle_cache,
1192                instrument_states,
1193                clock,
1194            );
1195        }
1196    }
1197}
1198
1199#[expect(clippy::too_many_arguments)]
1200fn handle_md_message(
1201    message: AxMdMessage,
1202    sender: &EventSender<DataEvent>,
1203    instruments: &Arc<AtomicMap<Ustr, InstrumentAny>>,
1204    symbol_data_types: &Arc<AtomicMap<String, SymbolDataTypes>>,
1205    book_sequences: &mut AHashMap<Ustr, u64>,
1206    candle_cache: &mut AHashMap<(Ustr, AxCandleWidth), AxMdCandle>,
1207    instrument_states: &mut AHashMap<Ustr, AxInstrumentState>,
1208    clock: &'static AtomicTime,
1209) {
1210    let ts_init = || -> UnixNanos { clock.get_time_ns() };
1211
1212    let instruments_snap = instruments.load();
1213    let sdt_snap = symbol_data_types.load();
1214
1215    match message {
1216        AxMdMessage::BookL1(book) => {
1217            let l1_subscribed = sdt_snap
1218                .get(book.s.as_str())
1219                .is_some_and(|e| e.quotes || e.book_level == Some(AxMarketDataLevel::Level1));
1220
1221            if !l1_subscribed {
1222                return;
1223            }
1224
1225            let Some(instrument) = instruments_snap.get(&book.s) else {
1226                log::error!(
1227                    "No instrument cached for symbol '{}' - cannot parse L1 book",
1228                    book.s
1229                );
1230                return;
1231            };
1232
1233            match parse_book_l1_quote(&book, instrument, ts_init()) {
1234                Ok(quote) => {
1235                    let _ = sender.send(DataEvent::Data(Data::Quote(quote)));
1236                }
1237                Err(e) => log::error!("Failed to parse L1 to QuoteTick: {e}"),
1238            }
1239        }
1240        AxMdMessage::BookL2(book) => {
1241            let symbol = book.s;
1242            let seq = book_sequences.entry(symbol).or_insert(0);
1243            *seq += 1;
1244            let sequence = *seq;
1245
1246            let Some(instrument) = instruments_snap.get(&symbol) else {
1247                log::error!("No instrument cached for symbol '{symbol}' - cannot parse L2 book");
1248                return;
1249            };
1250
1251            match parse_book_l2_deltas(&book, instrument, sequence, ts_init()) {
1252                Ok(deltas) => {
1253                    let _ = sender.send(DataEvent::Data(Data::BookDeltas(Box::new(deltas))));
1254                }
1255                Err(e) => log::error!("Failed to parse L2 to OrderBookDeltas: {e}"),
1256            }
1257
1258            let quotes_subscribed = sdt_snap
1259                .get(symbol.as_str())
1260                .is_some_and(|entry| entry.quotes);
1261
1262            if quotes_subscribed {
1263                match parse_book_l2_quote(&book, instrument, ts_init()) {
1264                    Ok(quote) => {
1265                        let _ = sender.send(DataEvent::Data(Data::Quote(quote)));
1266                    }
1267                    Err(e) => log::error!("Failed to parse L2 to QuoteTick: {e}"),
1268                }
1269            }
1270        }
1271        AxMdMessage::BookL3(book) => {
1272            let symbol = book.s;
1273            let seq = book_sequences.entry(symbol).or_insert(0);
1274            *seq += 1;
1275            let sequence = *seq;
1276
1277            let Some(instrument) = instruments_snap.get(&symbol) else {
1278                log::error!("No instrument cached for symbol '{symbol}' - cannot parse L3 book");
1279                return;
1280            };
1281
1282            match parse_book_l3_deltas(&book, instrument, sequence, ts_init()) {
1283                Ok(deltas) => {
1284                    let _ = sender.send(DataEvent::Data(Data::BookDeltas(Box::new(deltas))));
1285                }
1286                Err(e) => log::error!("Failed to parse L3 to OrderBookDeltas: {e}"),
1287            }
1288
1289            let quotes_subscribed = sdt_snap
1290                .get(symbol.as_str())
1291                .is_some_and(|entry| entry.quotes);
1292
1293            if quotes_subscribed {
1294                match parse_book_l3_quote(&book, instrument, ts_init()) {
1295                    Ok(quote) => {
1296                        let _ = sender.send(DataEvent::Data(Data::Quote(quote)));
1297                    }
1298                    Err(e) => log::error!("Failed to parse L3 to QuoteTick: {e}"),
1299                }
1300            }
1301        }
1302        AxMdMessage::Ticker(ticker) => {
1303            let Some(instrument) = instruments_snap.get(&ticker.s) else {
1304                log::debug!("No instrument cached for ticker symbol '{}'", ticker.s);
1305                return;
1306            };
1307
1308            let instrument_id = instrument.id();
1309            let price_precision = instrument.price_precision();
1310            let ts_event =
1311                ax_timestamp_stn_to_unix_nanos(ticker.ts, ticker.tn).unwrap_or_else(|_| ts_init());
1312            let ts_init = ts_init();
1313
1314            let mark_prices_subscribed = sdt_snap
1315                .get(ticker.s.as_str())
1316                .is_some_and(|e| e.mark_prices);
1317
1318            if mark_prices_subscribed && let Some(mark_price) = ticker.m {
1319                match Price::from_decimal_dp(mark_price, price_precision) {
1320                    Ok(price) => {
1321                        let update = MarkPriceUpdate::new(instrument_id, price, ts_event, ts_init);
1322                        let _ = sender.send(DataEvent::Data(Data::MarkPrice(update)));
1323                    }
1324                    Err(e) => {
1325                        log::error!("Failed to parse mark price for {}: {e}", ticker.s);
1326                    }
1327                }
1328            }
1329
1330            if let Some(state) = ticker.i {
1331                let status_subscribed = sdt_snap
1332                    .get(ticker.s.as_str())
1333                    .is_some_and(|e| e.instrument_status);
1334
1335                if status_subscribed {
1336                    let prev = instrument_states.insert(ticker.s, state);
1337                    if prev != Some(state) {
1338                        let action = MarketStatusAction::from(state);
1339                        let status = InstrumentStatus::new(
1340                            instrument_id,
1341                            action,
1342                            ts_event,
1343                            ts_init,
1344                            None,
1345                            None,
1346                            Some(state == AxInstrumentState::Open),
1347                            None,
1348                            None,
1349                        );
1350                        let _ = sender.send(DataEvent::InstrumentStatus(status));
1351                    }
1352                }
1353            }
1354        }
1355        AxMdMessage::Trade(trade) => {
1356            let trades_subscribed = sdt_snap.get(trade.s.as_str()).is_some_and(|e| e.trades);
1357
1358            if !trades_subscribed {
1359                return;
1360            }
1361
1362            let Some(instrument) = instruments_snap.get(&trade.s) else {
1363                log::error!(
1364                    "No instrument cached for symbol '{}' - cannot parse trade",
1365                    trade.s
1366                );
1367                return;
1368            };
1369
1370            match parse_trade_tick(&trade, instrument, ts_init()) {
1371                Ok(tick) => {
1372                    let _ = sender.send(DataEvent::Data(Data::Trade(tick)));
1373                }
1374                Err(e) => log::error!("Failed to parse trade to TradeTick: {e}"),
1375            }
1376        }
1377        AxMdMessage::Candle(candle) => {
1378            let cache_key = (candle.symbol, candle.width);
1379
1380            let closed_candle = if let Some(cached) = candle_cache.get(&cache_key) {
1381                if cached.ts == candle.ts {
1382                    None
1383                } else {
1384                    Some(cached.clone())
1385                }
1386            } else {
1387                None
1388            };
1389
1390            candle_cache.insert(cache_key, candle);
1391
1392            if let Some(closed) = closed_candle {
1393                let Some(instrument) = instruments_snap.get(&closed.symbol) else {
1394                    log::error!(
1395                        "No instrument cached for symbol '{}' - cannot parse candle",
1396                        closed.symbol
1397                    );
1398                    return;
1399                };
1400
1401                match parse_candle_bar(&closed, instrument, ts_init()) {
1402                    Ok(bar) => {
1403                        let _ = sender.send(DataEvent::Data(Data::Bar(bar)));
1404                    }
1405                    Err(e) => log::error!("Failed to parse candle to Bar: {e}"),
1406                }
1407            }
1408        }
1409        AxMdMessage::Heartbeat(_) => {
1410            log::trace!("Received heartbeat");
1411        }
1412        AxMdMessage::SubscriptionResponse(_) => {}
1413        AxMdMessage::Error(error) => {
1414            log::warn!("WebSocket error: {}", error.message);
1415        }
1416    }
1417}
1418
1419#[cfg(test)]
1420mod tests {
1421    use std::sync::Arc;
1422
1423    use ahash::{AHashMap, AHashSet};
1424    use nautilus_model::{
1425        data::InstrumentStatus,
1426        enums::AssetClass,
1427        identifiers::{InstrumentId, Symbol},
1428        instruments::PerpetualContract,
1429        types::{Currency, Price, Quantity},
1430    };
1431    use parking_lot::Mutex;
1432    use rstest::rstest;
1433    use rust_decimal::Decimal;
1434    use rust_decimal_macros::dec;
1435    use ustr::Ustr;
1436
1437    use super::*;
1438    use crate::websocket::{
1439        data::client::SymbolDataTypes,
1440        messages::{AxBookLevel, AxMdBookL2, AxMdMessage, AxMdTicker},
1441    };
1442
1443    #[rstest]
1444    fn test_drain_status_invalidations_removes_cached_state() {
1445        let invalidations = Arc::new(Mutex::new(AHashSet::new()));
1446        let mut states = AHashMap::new();
1447        let sym = Ustr::from("EURUSD-PERP");
1448
1449        states.insert(sym, AxInstrumentState::Open);
1450        invalidations.lock().insert(sym);
1451
1452        drain_status_invalidations(&invalidations, &mut states);
1453
1454        assert!(!states.contains_key(&sym));
1455        assert!(invalidations.lock().is_empty());
1456    }
1457
1458    #[rstest]
1459    fn test_drain_status_invalidations_no_op_when_empty() {
1460        let invalidations = Arc::new(Mutex::new(AHashSet::new()));
1461        let mut states = AHashMap::new();
1462        let sym = Ustr::from("EURUSD-PERP");
1463        states.insert(sym, AxInstrumentState::Open);
1464
1465        drain_status_invalidations(&invalidations, &mut states);
1466
1467        assert!(states.contains_key(&sym));
1468    }
1469
1470    fn ticker_test_instrument() -> InstrumentAny {
1471        let symbol = Symbol::new("EURUSD-PERP");
1472        let instrument = PerpetualContract::builder()
1473            .instrument_id(InstrumentId::new(symbol, *crate::common::consts::AX_VENUE))
1474            .raw_symbol(symbol)
1475            .underlying(Ustr::from("EURUSD"))
1476            .asset_class(AssetClass::FX)
1477            .quote_currency(Currency::USD())
1478            .settlement_currency(Currency::USD())
1479            .is_inverse(false)
1480            .price_precision(4)
1481            .size_precision(0)
1482            .price_increment(Price::from("0.0001"))
1483            .size_increment(Quantity::from("1"))
1484            .margin_init(Decimal::new(1, 2))
1485            .margin_maint(Decimal::new(5, 3))
1486            .maker_fee(Decimal::new(2, 4))
1487            .taker_fee(Decimal::new(5, 4))
1488            .ts_event(UnixNanos::default())
1489            .ts_init(UnixNanos::default())
1490            .build()
1491            .unwrap();
1492        InstrumentAny::PerpetualContract(instrument)
1493    }
1494
1495    fn ticker_message(state: AxInstrumentState) -> AxMdTicker {
1496        AxMdTicker {
1497            ts: 1_700_000_000,
1498            tn: 0,
1499            s: Ustr::from("EURUSD-PERP"),
1500            p: rust_decimal::Decimal::ZERO,
1501            q: 0,
1502            o: rust_decimal::Decimal::ZERO,
1503            l: rust_decimal::Decimal::ZERO,
1504            h: rust_decimal::Decimal::ZERO,
1505            v: 0,
1506            oi: None,
1507            m: None,
1508            i: Some(state),
1509            pl: None,
1510            pu: None,
1511            lsp: None,
1512        }
1513    }
1514
1515    fn collect_instrument_statuses(
1516        rx: &mut tokio::sync::mpsc::UnboundedReceiver<DataEvent>,
1517    ) -> Vec<InstrumentStatus> {
1518        let mut statuses = Vec::new();
1519
1520        while let Ok(event) = rx.try_recv() {
1521            if let DataEvent::InstrumentStatus(status) = event {
1522                statuses.push(status);
1523            }
1524        }
1525        statuses
1526    }
1527
1528    #[rstest]
1529    fn test_ticker_instrument_status_emitted_once_when_state_unchanged() {
1530        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1531        let instruments = Arc::new(AtomicMap::new());
1532        instruments.insert(Ustr::from("EURUSD-PERP"), ticker_test_instrument());
1533
1534        let sdt = Arc::new(AtomicMap::new());
1535        sdt.insert(
1536            "EURUSD-PERP".to_string(),
1537            SymbolDataTypes {
1538                quotes: false,
1539                trades: false,
1540                mark_prices: false,
1541                instrument_status: true,
1542                book_level: None,
1543            },
1544        );
1545
1546        let mut book_sequences = AHashMap::new();
1547        let mut candle_cache = AHashMap::new();
1548        let mut instrument_states = AHashMap::new();
1549        let clock = get_atomic_clock_realtime();
1550
1551        let msg = AxMdMessage::Ticker(ticker_message(AxInstrumentState::Open));
1552        handle_md_message(
1553            msg.clone(),
1554            &tx.clone().into(),
1555            &instruments,
1556            &sdt,
1557            &mut book_sequences,
1558            &mut candle_cache,
1559            &mut instrument_states,
1560            clock,
1561        );
1562
1563        // Same state repeated: second call should not emit a second InstrumentStatus
1564        handle_md_message(
1565            msg,
1566            &tx.into(),
1567            &instruments,
1568            &sdt,
1569            &mut book_sequences,
1570            &mut candle_cache,
1571            &mut instrument_states,
1572            clock,
1573        );
1574
1575        let statuses = collect_instrument_statuses(&mut rx);
1576        assert_eq!(
1577            statuses.len(),
1578            1,
1579            "expected a single emission, found {statuses:?}"
1580        );
1581        assert_eq!(statuses[0].is_trading, Some(true));
1582    }
1583
1584    #[rstest]
1585    fn test_ticker_instrument_status_emitted_on_transition() {
1586        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1587        let instruments = Arc::new(AtomicMap::new());
1588        instruments.insert(Ustr::from("EURUSD-PERP"), ticker_test_instrument());
1589
1590        let sdt = Arc::new(AtomicMap::new());
1591        sdt.insert(
1592            "EURUSD-PERP".to_string(),
1593            SymbolDataTypes {
1594                quotes: false,
1595                trades: false,
1596                mark_prices: false,
1597                instrument_status: true,
1598                book_level: None,
1599            },
1600        );
1601
1602        let mut book_sequences = AHashMap::new();
1603        let mut candle_cache = AHashMap::new();
1604        let mut instrument_states = AHashMap::new();
1605        let clock = get_atomic_clock_realtime();
1606
1607        handle_md_message(
1608            AxMdMessage::Ticker(ticker_message(AxInstrumentState::Open)),
1609            &tx.clone().into(),
1610            &instruments,
1611            &sdt,
1612            &mut book_sequences,
1613            &mut candle_cache,
1614            &mut instrument_states,
1615            clock,
1616        );
1617        handle_md_message(
1618            AxMdMessage::Ticker(ticker_message(AxInstrumentState::Closed)),
1619            &tx.into(),
1620            &instruments,
1621            &sdt,
1622            &mut book_sequences,
1623            &mut candle_cache,
1624            &mut instrument_states,
1625            clock,
1626        );
1627
1628        let statuses = collect_instrument_statuses(&mut rx);
1629        assert_eq!(statuses.len(), 2, "expected one emission per transition");
1630        assert_eq!(statuses[0].is_trading, Some(true));
1631        assert_eq!(statuses[1].is_trading, Some(false));
1632    }
1633
1634    #[rstest]
1635    fn test_ticker_instrument_status_skipped_when_not_subscribed() {
1636        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1637        let instruments = Arc::new(AtomicMap::new());
1638        instruments.insert(Ustr::from("EURUSD-PERP"), ticker_test_instrument());
1639
1640        let sdt = Arc::new(AtomicMap::new());
1641        sdt.insert(
1642            "EURUSD-PERP".to_string(),
1643            SymbolDataTypes {
1644                quotes: false,
1645                trades: false,
1646                mark_prices: false,
1647                instrument_status: false,
1648                book_level: None,
1649            },
1650        );
1651
1652        let mut book_sequences = AHashMap::new();
1653        let mut candle_cache = AHashMap::new();
1654        let mut instrument_states = AHashMap::new();
1655        let clock = get_atomic_clock_realtime();
1656
1657        handle_md_message(
1658            AxMdMessage::Ticker(ticker_message(AxInstrumentState::Open)),
1659            &tx.into(),
1660            &instruments,
1661            &sdt,
1662            &mut book_sequences,
1663            &mut candle_cache,
1664            &mut instrument_states,
1665            clock,
1666        );
1667
1668        let statuses = collect_instrument_statuses(&mut rx);
1669        assert!(statuses.is_empty());
1670    }
1671
1672    #[rstest]
1673    fn test_l2_book_emits_quote_when_quotes_subscribed() {
1674        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel();
1675        let instruments = Arc::new(AtomicMap::new());
1676        instruments.insert(Ustr::from("EURUSD-PERP"), ticker_test_instrument());
1677
1678        let sdt = Arc::new(AtomicMap::new());
1679        sdt.insert(
1680            "EURUSD-PERP".to_string(),
1681            SymbolDataTypes {
1682                quotes: true,
1683                book_level: Some(AxMarketDataLevel::Level2),
1684                ..Default::default()
1685            },
1686        );
1687
1688        let mut book_sequences = AHashMap::new();
1689        let mut candle_cache = AHashMap::new();
1690        let mut instrument_states = AHashMap::new();
1691        let clock = get_atomic_clock_realtime();
1692        let message = AxMdMessage::BookL2(AxMdBookL2 {
1693            ts: 1_700_000_000,
1694            tn: 123,
1695            s: Ustr::from("EURUSD-PERP"),
1696            b: vec![AxBookLevel {
1697                p: dec!(1.1441),
1698                q: 100,
1699            }],
1700            a: vec![AxBookLevel {
1701                p: dec!(1.1448),
1702                q: 200,
1703            }],
1704            st: true,
1705        });
1706
1707        handle_md_message(
1708            message,
1709            &tx.into(),
1710            &instruments,
1711            &sdt,
1712            &mut book_sequences,
1713            &mut candle_cache,
1714            &mut instrument_states,
1715            clock,
1716        );
1717
1718        let events = std::iter::from_fn(|| rx.try_recv().ok()).collect::<Vec<_>>();
1719        let quote = events.iter().find_map(|event| match event {
1720            DataEvent::Data(Data::Quote(quote)) => Some(quote),
1721            _ => None,
1722        });
1723
1724        assert_eq!(
1725            quote.map(|quote| quote.bid_price),
1726            Some(Price::from("1.1441"))
1727        );
1728        assert_eq!(
1729            quote.map(|quote| quote.ask_price),
1730            Some(Price::from("1.1448"))
1731        );
1732    }
1733}