1use 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#[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 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 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 #[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 #[must_use]
177 pub fn instruments(&self) -> &Arc<AtomicMap<Ustr, InstrumentAny>> {
178 &self.instruments
179 }
180
181 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 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 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 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 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 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 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 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}