Skip to main content

nautilus_kraken/data/
futures.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//! Kraken Futures data client implementation.
17
18use std::{
19    future::Future,
20    sync::{
21        Arc,
22        atomic::{AtomicBool, AtomicU64, Ordering},
23    },
24    time::Duration,
25};
26
27use ahash::AHashMap;
28use anyhow::Context;
29use async_trait::async_trait;
30use nautilus_common::{
31    clients::DataClient,
32    live::{get_data_event_sender, sender::EventSender},
33    messages::{
34        DataEvent,
35        data::{
36            BarsResponse, BookResponse, DataResponse, FundingRatesResponse, InstrumentResponse,
37            InstrumentsResponse, RequestBars, RequestBookSnapshot, RequestFundingRates,
38            RequestInstrument, RequestInstruments, RequestTrades, SubscribeBars,
39            SubscribeBookDeltas, SubscribeFundingRates, SubscribeIndexPrices, SubscribeInstrument,
40            SubscribeInstrumentStatus, SubscribeInstruments, SubscribeMarkPrices, SubscribeQuotes,
41            SubscribeTrades, TradesResponse, UnsubscribeBars, UnsubscribeBookDeltas,
42            UnsubscribeFundingRates, UnsubscribeIndexPrices, UnsubscribeInstrumentStatus,
43            UnsubscribeMarkPrices, UnsubscribeQuotes, UnsubscribeTrades,
44        },
45    },
46};
47use nautilus_core::{
48    AtomicMap, AtomicSet,
49    datetime::datetime_to_unix_nanos,
50    nanos::UnixNanos,
51    time::{AtomicTime, get_atomic_clock_realtime},
52};
53use nautilus_live::{SocketControl, task::TaskGroup};
54use nautilus_model::{
55    data::{Data, OrderBookDeltas, QuoteTick},
56    enums::BookType,
57    identifiers::{ClientId, InstrumentId, Venue},
58    instruments::{Instrument, InstrumentAny},
59    orderbook::OrderBook,
60};
61use rust_decimal_macros::dec;
62use tokio_util::sync::CancellationToken;
63
64use crate::{
65    common::{consts::KRAKEN_VENUE, lookup_instrument_in_snapshot},
66    config::KrakenDataClientConfig,
67    http::{
68        KrakenFuturesHttpClient, futures::client::KRAKEN_FUTURES_DEFAULT_RATE_LIMIT_PER_SECOND,
69    },
70    websocket::futures::{
71        client::KrakenFuturesWebSocketClient,
72        messages::KrakenFuturesWsMessage,
73        parse::{
74            parse_futures_ws_book_delta, parse_futures_ws_book_snapshot_deltas,
75            parse_futures_ws_funding_rate, parse_futures_ws_index_price,
76            parse_futures_ws_mark_price, parse_futures_ws_trade_tick,
77        },
78    },
79};
80
81/// Kraken Futures data client.
82///
83/// Provides real-time market data from Kraken Futures markets.
84#[allow(dead_code)]
85#[derive(Debug)]
86pub struct KrakenFuturesDataClient {
87    clock: &'static AtomicTime,
88    client_id: ClientId,
89    config: KrakenDataClientConfig,
90    http: KrakenFuturesHttpClient,
91    ws: KrakenFuturesWebSocketClient,
92    is_connected: AtomicBool,
93    cancellation_token: CancellationToken,
94    session_tasks: TaskGroup,
95    command_tasks: TaskGroup,
96    instruments: Arc<AtomicMap<InstrumentId, InstrumentAny>>,
97    quote_instruments: Arc<AtomicSet<InstrumentId>>,
98    book_instruments: Arc<AtomicSet<InstrumentId>>,
99    data_sender: EventSender<DataEvent>,
100}
101
102impl KrakenFuturesDataClient {
103    /// Creates a new [`KrakenFuturesDataClient`] instance.
104    pub fn new(client_id: ClientId, config: KrakenDataClientConfig) -> anyhow::Result<Self> {
105        let session_tasks = TaskGroup::new();
106        let cancellation_token = session_tasks.cancellation_token();
107        let command_tasks = TaskGroup::new();
108        let proxy_url = config
109            .proxy_url
110            .as_ref()
111            .map(|value| value.expose_secret().to_owned());
112
113        let http = KrakenFuturesHttpClient::new(
114            config.environment,
115            config.base_url.clone(),
116            config.timeout_secs,
117            None,
118            None,
119            None,
120            proxy_url.clone(),
121            config
122                .max_requests_per_second
123                .unwrap_or(KRAKEN_FUTURES_DEFAULT_RATE_LIMIT_PER_SECOND),
124        )?;
125
126        let ws = KrakenFuturesWebSocketClient::with_credentials(
127            config.ws_public_url(),
128            config.heartbeat_interval_secs,
129            None,
130            None,
131            config.transport_backend,
132            proxy_url,
133        )
134        .with_socket_control(SocketControl::new(
135            client_id,
136            Some(*KRAKEN_VENUE),
137            "kraken-futures-data-streams",
138        ));
139
140        Ok(Self {
141            clock: get_atomic_clock_realtime(),
142            client_id,
143            config,
144            http,
145            ws,
146            is_connected: AtomicBool::new(false),
147            cancellation_token,
148            session_tasks,
149            command_tasks,
150            instruments: Arc::new(AtomicMap::new()),
151            quote_instruments: Arc::new(AtomicSet::new()),
152            book_instruments: Arc::new(AtomicSet::new()),
153            data_sender: get_data_event_sender(),
154        })
155    }
156
157    /// Returns the cached instruments.
158    #[must_use]
159    pub fn instruments(&self) -> Vec<InstrumentAny> {
160        self.instruments.load().values().cloned().collect()
161    }
162
163    /// Returns a cached instrument by ID.
164    #[must_use]
165    pub fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
166        self.instruments.load().get(instrument_id).cloned()
167    }
168
169    async fn load_instruments(&self) -> anyhow::Result<Vec<InstrumentAny>> {
170        let instruments = self
171            .http
172            .request_instruments()
173            .await
174            .context("Failed to load futures instruments")?;
175
176        self.instruments.rcu(|m| {
177            for instrument in &instruments {
178                m.insert(instrument.id(), instrument.clone());
179            }
180        });
181
182        self.http.cache_instruments(&instruments);
183
184        log::debug!(
185            "Loaded instruments: client_id={}, count={}",
186            self.client_id,
187            instruments.len()
188        );
189
190        Ok(instruments)
191    }
192
193    fn spawn_ws<F>(&self, fut: F, context: &'static str)
194    where
195        F: Future<Output = anyhow::Result<()>> + Send + 'static,
196    {
197        let future = async move {
198            if let Err(e) = fut.await {
199                log::error!("{context}: {e:?}");
200            }
201        };
202
203        if let Err(e) = self.command_tasks.spawn(future) {
204            log::warn!("Skipping Kraken Futures {context} after shutdown began: {e}");
205        }
206    }
207
208    fn spawn_command<F>(&self, future: F)
209    where
210        F: Future<Output = ()> + Send + 'static,
211    {
212        if let Err(e) = self.command_tasks.spawn(future) {
213            log::warn!("Skipping Kraken Futures data command after shutdown began: {e}");
214        }
215    }
216
217    async fn finish_tasks(&self) -> anyhow::Result<()> {
218        let (session_result, command_result) = tokio::join!(
219            self.session_tasks
220                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
221            self.command_tasks
222                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
223        );
224        session_result.context("failed to finish Kraken Futures data session tasks")?;
225        command_result.context("failed to finish Kraken Futures data command tasks")?;
226        Ok(())
227    }
228
229    async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
230        if !self.session_tasks.is_open() || !self.command_tasks.is_open() {
231            self.session_tasks.begin_shutdown();
232            self.command_tasks.begin_shutdown();
233            let _ = self.ws.close().await;
234            self.finish_tasks().await?;
235            self.session_tasks
236                .start_generation()
237                .context("failed to start Kraken Futures data session task generation")?;
238            self.command_tasks
239                .start_generation()
240                .context("failed to start Kraken Futures data command task generation")?;
241            self.cancellation_token = self.session_tasks.cancellation_token();
242            self.ws = KrakenFuturesWebSocketClient::with_credentials(
243                self.config.ws_public_url(),
244                self.config.heartbeat_interval_secs,
245                None,
246                None,
247                self.config.transport_backend,
248                self.config
249                    .proxy_url
250                    .as_ref()
251                    .map(|value| value.expose_secret().to_owned()),
252            )
253            .with_socket_control(SocketControl::new(
254                self.client_id,
255                Some(*KRAKEN_VENUE),
256                "kraken-futures-data-streams",
257            ));
258        }
259        Ok(())
260    }
261
262    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
263        self.session_tasks.begin_shutdown();
264        self.command_tasks.begin_shutdown();
265        let ws_result = self.ws.close().await;
266        let tasks_result = self.finish_tasks().await;
267        self.is_connected.store(false, Ordering::Release);
268        tasks_result?;
269        Ok(ws_result?)
270    }
271
272    fn spawn_message_handler(&mut self) -> anyhow::Result<()> {
273        let mut rx = self
274            .ws
275            .take_output_rx()
276            .context("Failed to take futures WebSocket output receiver")?;
277        let data_sender = self.data_sender.clone();
278        let instruments = self.instruments.clone();
279        let quote_instruments = self.quote_instruments.clone();
280        let book_instruments = self.book_instruments.clone();
281        let book_sequence = Arc::new(AtomicU64::new(0));
282        let cancellation_token = self.cancellation_token.clone();
283        let clock = self.clock;
284
285        let future = async move {
286            let mut order_books: AHashMap<InstrumentId, OrderBook> = AHashMap::new();
287            let mut last_quotes: AHashMap<InstrumentId, QuoteTick> = AHashMap::new();
288
289            loop {
290                tokio::select! {
291                    () = cancellation_token.cancelled() => {
292                        log::debug!("Futures message handler cancelled");
293                        break;
294                    }
295                    msg = rx.recv() => {
296                        match msg {
297                            Some(ws_msg) => {
298                                Self::handle_ws_message(
299                                    ws_msg,
300                                    &data_sender,
301                                    &instruments,
302                                    &quote_instruments,
303                                    &book_instruments,
304                                    &mut order_books,
305                                    &mut last_quotes,
306                                    &book_sequence,
307                                    clock,
308                                );
309                            }
310                            None => {
311                                log::debug!("Futures WebSocket stream ended");
312                                break;
313                            }
314                        }
315                    }
316                }
317            }
318        };
319
320        self.session_tasks
321            .spawn(future)
322            .context("failed to register Kraken Futures message handler")
323    }
324
325    #[expect(clippy::too_many_arguments)]
326    fn handle_ws_message(
327        msg: KrakenFuturesWsMessage,
328        sender: &EventSender<DataEvent>,
329        instruments: &Arc<AtomicMap<InstrumentId, InstrumentAny>>,
330        quote_instruments: &Arc<AtomicSet<InstrumentId>>,
331        book_instruments: &Arc<AtomicSet<InstrumentId>>,
332        order_books: &mut AHashMap<InstrumentId, OrderBook>,
333        last_quotes: &mut AHashMap<InstrumentId, QuoteTick>,
334        book_sequence: &Arc<AtomicU64>,
335        clock: &'static AtomicTime,
336    ) {
337        let ts_init = clock.get_time_ns();
338
339        match msg {
340            KrakenFuturesWsMessage::Ticker(ticker) => {
341                let instruments = instruments.load();
342                let Some(instrument) =
343                    lookup_instrument_in_snapshot(&instruments, ticker.product_id.as_str())
344                else {
345                    log::warn!("No instrument for product_id: {}", ticker.product_id);
346                    return;
347                };
348
349                if let Some(mark) = parse_futures_ws_mark_price(&ticker, instrument, ts_init)
350                    && let Err(e) = sender.send(DataEvent::Data(Data::MarkPrice(mark)))
351                {
352                    log::error!("Failed to send mark price: {e}");
353                }
354
355                if let Some(index) = parse_futures_ws_index_price(&ticker, instrument, ts_init)
356                    && let Err(e) = sender.send(DataEvent::Data(Data::IndexPrice(index)))
357                {
358                    log::error!("Failed to send index price: {e}");
359                }
360
361                if let Some(funding) = parse_futures_ws_funding_rate(&ticker, instrument, ts_init)
362                    && let Err(e) = sender.send(DataEvent::FundingRate(funding))
363                {
364                    log::error!("Failed to send funding rate: {e}");
365                }
366            }
367            KrakenFuturesWsMessage::Trade(trade) => {
368                let instruments = instruments.load();
369                let Some(instrument) =
370                    lookup_instrument_in_snapshot(&instruments, trade.product_id.as_str())
371                else {
372                    log::warn!("No instrument for product_id: {}", trade.product_id);
373                    return;
374                };
375
376                match parse_futures_ws_trade_tick(&trade, instrument, ts_init) {
377                    Ok(tick) => {
378                        if let Err(e) = sender.send(DataEvent::Data(Data::Trade(tick))) {
379                            log::error!("Failed to send trade: {e}");
380                        }
381                    }
382                    Err(e) => log::error!("Failed to parse futures trade tick: {e}"),
383                }
384            }
385            KrakenFuturesWsMessage::BookSnapshot(snapshot) => {
386                let instruments = instruments.load();
387                let Some(instrument) =
388                    lookup_instrument_in_snapshot(&instruments, snapshot.product_id.as_str())
389                else {
390                    log::warn!("No instrument for product_id: {}", snapshot.product_id);
391                    return;
392                };
393                let instrument_id = instrument.id();
394                let sequence = book_sequence.load(Ordering::Relaxed);
395
396                match parse_futures_ws_book_snapshot_deltas(
397                    &snapshot, instrument, sequence, ts_init,
398                ) {
399                    Ok(delta_vec) => {
400                        if delta_vec.is_empty() {
401                            return;
402                        }
403                        book_sequence.fetch_add(delta_vec.len() as u64, Ordering::Relaxed);
404                        let deltas = OrderBookDeltas::new(instrument_id, delta_vec);
405
406                        let has_quote_sub = quote_instruments.contains(&instrument_id);
407
408                        if has_quote_sub {
409                            let book = order_books
410                                .entry(instrument_id)
411                                .or_insert_with(|| OrderBook::new(instrument_id, BookType::L2_MBP));
412
413                            if let Err(e) = book.apply_deltas(&deltas) {
414                                log::error!("Failed to apply snapshot deltas to order book: {e}");
415                            } else {
416                                Self::maybe_emit_quote(
417                                    book,
418                                    instrument_id,
419                                    last_quotes,
420                                    ts_init,
421                                    sender,
422                                );
423                            }
424                        }
425
426                        let has_book_sub = book_instruments.contains(&instrument_id);
427
428                        if has_book_sub
429                            && let Err(e) =
430                                sender.send(DataEvent::Data(Data::BookDeltas(Box::new(deltas))))
431                        {
432                            log::error!("Failed to send book snapshot deltas: {e}");
433                        }
434                    }
435                    Err(e) => log::error!("Failed to parse book snapshot: {e}"),
436                }
437            }
438            KrakenFuturesWsMessage::BookDelta(delta) => {
439                let instruments = instruments.load();
440                let Some(instrument) =
441                    lookup_instrument_in_snapshot(&instruments, delta.product_id.as_str())
442                else {
443                    log::warn!("No instrument for product_id: {}", delta.product_id);
444                    return;
445                };
446                let instrument_id = instrument.id();
447                let sequence = book_sequence.fetch_add(1, Ordering::Relaxed);
448                match parse_futures_ws_book_delta(&delta, instrument, sequence, ts_init) {
449                    Ok(book_delta) => {
450                        let deltas = OrderBookDeltas::new(instrument_id, vec![book_delta]);
451
452                        let has_quote_sub = quote_instruments.contains(&instrument_id);
453
454                        if has_quote_sub && let Some(book) = order_books.get_mut(&instrument_id) {
455                            if let Err(e) = book.apply_deltas(&deltas) {
456                                log::error!("Failed to apply delta to order book: {e}");
457                            } else {
458                                Self::maybe_emit_quote(
459                                    book,
460                                    instrument_id,
461                                    last_quotes,
462                                    ts_init,
463                                    sender,
464                                );
465                            }
466                        }
467
468                        let has_book_sub = book_instruments.contains(&instrument_id);
469
470                        if has_book_sub
471                            && let Err(e) =
472                                sender.send(DataEvent::Data(Data::BookDeltas(Box::new(deltas))))
473                        {
474                            log::error!("Failed to send book delta: {e}");
475                        }
476                    }
477                    Err(e) => log::error!("Failed to parse book delta: {e}"),
478                }
479            }
480            KrakenFuturesWsMessage::Reconnected => {
481                log::info!("Futures WebSocket reconnected");
482            }
483            KrakenFuturesWsMessage::OpenOrdersCancel(_)
484            | KrakenFuturesWsMessage::OpenOrdersDelta(_)
485            | KrakenFuturesWsMessage::FillsDelta(_)
486            | KrakenFuturesWsMessage::Challenge(_) => {}
487        }
488    }
489
490    fn maybe_emit_quote(
491        book: &OrderBook,
492        instrument_id: InstrumentId,
493        last_quotes: &mut AHashMap<InstrumentId, QuoteTick>,
494        ts_init: UnixNanos,
495        sender: &EventSender<DataEvent>,
496    ) {
497        let (Some(bid_price), Some(ask_price)) = (book.best_bid_price(), book.best_ask_price())
498        else {
499            return;
500        };
501        let (Some(bid_size), Some(ask_size)) = (book.best_bid_size(), book.best_ask_size()) else {
502            return;
503        };
504
505        let bid = bid_price.as_decimal();
506        let ask = ask_price.as_decimal();
507        if bid > dec!(0) && (ask - bid) / bid > dec!(0.25) {
508            log::debug!("Filtered quote with wide spread: bid={bid}, ask={ask}");
509            return;
510        }
511
512        let quote = QuoteTick::new(
513            instrument_id,
514            bid_price,
515            ask_price,
516            bid_size,
517            ask_size,
518            ts_init,
519            ts_init,
520        );
521
522        if matches!(last_quotes.get(&instrument_id), Some(prev) if *prev == quote) {
523            return;
524        }
525
526        last_quotes.insert(instrument_id, quote);
527
528        if let Err(e) = sender.send(DataEvent::Data(Data::Quote(quote))) {
529            log::error!("Failed to send quote: {e}");
530        }
531    }
532}
533
534#[async_trait(?Send)]
535impl DataClient for KrakenFuturesDataClient {
536    fn client_id(&self) -> ClientId {
537        self.client_id
538    }
539
540    fn venue(&self) -> Option<Venue> {
541        Some(*KRAKEN_VENUE)
542    }
543
544    fn start(&mut self) -> anyhow::Result<()> {
545        log::info!(
546            "Starting Futures data client: client_id={}, environment={:?}",
547            self.client_id,
548            self.config.environment
549        );
550        Ok(())
551    }
552
553    fn stop(&mut self) -> anyhow::Result<()> {
554        log::info!("Stopping Futures data client: {}", self.client_id);
555        self.session_tasks.begin_shutdown();
556        self.command_tasks.begin_shutdown();
557        self.ws.begin_shutdown();
558        self.is_connected.store(false, Ordering::Relaxed);
559        Ok(())
560    }
561
562    fn reset(&mut self) -> anyhow::Result<()> {
563        log::info!("Resetting Futures data client: {}", self.client_id);
564        self.session_tasks.begin_shutdown();
565        self.command_tasks.begin_shutdown();
566        self.ws.begin_shutdown();
567        self.is_connected.store(false, Ordering::Relaxed);
568
569        self.instruments.store(ahash::AHashMap::new());
570
571        self.quote_instruments.store(ahash::AHashSet::new());
572
573        Ok(())
574    }
575
576    fn dispose(&mut self) -> anyhow::Result<()> {
577        log::debug!("Disposing Futures data client: {}", self.client_id);
578        self.stop()
579    }
580
581    fn is_connected(&self) -> bool {
582        self.is_connected.load(Ordering::SeqCst)
583    }
584
585    fn is_disconnected(&self) -> bool {
586        !self.is_connected()
587    }
588
589    async fn connect(&mut self) -> anyhow::Result<()> {
590        if self.is_connected() && self.session_tasks.is_open() && self.command_tasks.is_open() {
591            return Ok(());
592        }
593
594        self.prepare_task_groups().await?;
595
596        let instruments = self.load_instruments().await?;
597
598        let session_result = async {
599            self.ws
600                .connect()
601                .await
602                .context("Failed to connect futures WebSocket")?;
603            self.ws
604                .wait_until_active(10.0)
605                .await
606                .context("Futures WebSocket failed to become active")?;
607
608            self.spawn_message_handler()?;
609
610            Ok::<(), anyhow::Error>(())
611        }
612        .await;
613
614        if let Err(e) = session_result {
615            if let Err(teardown_error) = self.teardown_partial_connect().await {
616                return Err(e.context(format!(
617                    "Kraken Futures data startup teardown failed: {teardown_error}"
618                )));
619            }
620            return Err(e);
621        }
622
623        for instrument in instruments {
624            if let Err(e) = self.data_sender.send(DataEvent::Instrument(instrument)) {
625                log::error!("Failed to send instrument: {e}");
626            }
627        }
628
629        self.is_connected.store(true, Ordering::Release);
630        log::info!(
631            "Connected: client_id={}, product_type=Futures",
632            self.client_id
633        );
634        Ok(())
635    }
636
637    async fn disconnect(&mut self) -> anyhow::Result<()> {
638        self.teardown_partial_connect().await?;
639
640        self.quote_instruments.store(ahash::AHashSet::new());
641        self.is_connected.store(false, Ordering::Relaxed);
642
643        log::info!("Disconnected: client_id={}", self.client_id);
644        Ok(())
645    }
646
647    fn subscribe_instruments(&mut self, _cmd: SubscribeInstruments) -> anyhow::Result<()> {
648        log::debug!("subscribe_instruments: Kraken instruments are fetched via HTTP on connect");
649        Ok(())
650    }
651
652    fn subscribe_instrument(&mut self, _cmd: SubscribeInstrument) -> anyhow::Result<()> {
653        log::debug!("subscribe_instrument: Kraken instruments are fetched via HTTP on connect");
654        Ok(())
655    }
656
657    fn subscribe_book_deltas(&mut self, cmd: SubscribeBookDeltas) -> anyhow::Result<()> {
658        let instrument_id = cmd.instrument_id;
659        let depth = cmd.depth;
660
661        if cmd.book_type != BookType::L2_MBP {
662            log::warn!(
663                "Book type {:?} not supported by Kraken, skipping subscription",
664                cmd.book_type
665            );
666            return Ok(());
667        }
668
669        self.book_instruments.insert(instrument_id);
670
671        let ws = self.ws.clone();
672        self.spawn_ws(
673            async move {
674                ws.subscribe_book(instrument_id, depth.map(|d| d.get() as u32))
675                    .await
676                    .map_err(|e| anyhow::anyhow!("{e}"))
677            },
678            "subscribe book",
679        );
680
681        Ok(())
682    }
683
684    fn subscribe_quotes(&mut self, cmd: SubscribeQuotes) -> anyhow::Result<()> {
685        let instrument_id = cmd.instrument_id;
686        let ws = self.ws.clone();
687
688        self.quote_instruments.insert(instrument_id);
689
690        self.spawn_ws(
691            async move {
692                ws.subscribe_quotes(instrument_id)
693                    .await
694                    .map_err(|e| anyhow::anyhow!("{e}"))
695            },
696            "subscribe quotes",
697        );
698
699        Ok(())
700    }
701
702    fn subscribe_trades(&mut self, cmd: SubscribeTrades) -> anyhow::Result<()> {
703        let instrument_id = cmd.instrument_id;
704        let ws = self.ws.clone();
705
706        self.spawn_ws(
707            async move {
708                ws.subscribe_trades(instrument_id)
709                    .await
710                    .map_err(|e| anyhow::anyhow!("{e}"))
711            },
712            "subscribe trades",
713        );
714
715        Ok(())
716    }
717
718    fn subscribe_mark_prices(&mut self, cmd: SubscribeMarkPrices) -> anyhow::Result<()> {
719        let instrument_id = cmd.instrument_id;
720        let ws = self.ws.clone();
721
722        self.spawn_ws(
723            async move {
724                ws.subscribe_mark_price(instrument_id)
725                    .await
726                    .map_err(|e| anyhow::anyhow!("{e}"))
727            },
728            "subscribe mark price",
729        );
730
731        Ok(())
732    }
733
734    fn subscribe_index_prices(&mut self, cmd: SubscribeIndexPrices) -> anyhow::Result<()> {
735        let instrument_id = cmd.instrument_id;
736        let ws = self.ws.clone();
737
738        self.spawn_ws(
739            async move {
740                ws.subscribe_index_price(instrument_id)
741                    .await
742                    .map_err(|e| anyhow::anyhow!("{e}"))
743            },
744            "subscribe index price",
745        );
746
747        Ok(())
748    }
749
750    fn subscribe_funding_rates(&mut self, cmd: SubscribeFundingRates) -> anyhow::Result<()> {
751        let instrument_id = cmd.instrument_id;
752        let ws = self.ws.clone();
753
754        self.spawn_ws(
755            async move {
756                ws.subscribe_funding_rate(instrument_id)
757                    .await
758                    .map_err(|e| anyhow::anyhow!("{e}"))
759            },
760            "subscribe funding rate",
761        );
762
763        Ok(())
764    }
765
766    fn subscribe_bars(&mut self, cmd: SubscribeBars) -> anyhow::Result<()> {
767        log::warn!(
768            "Cannot subscribe to {} bars: Kraken Futures does not support EXTERNAL bar streaming",
769            cmd.bar_type
770        );
771        Ok(())
772    }
773
774    fn subscribe_instrument_status(
775        &mut self,
776        cmd: SubscribeInstrumentStatus,
777    ) -> anyhow::Result<()> {
778        log::debug!(
779            "subscribe_instrument_status: {} (status changes detected via periodic instrument polling)",
780            cmd.instrument_id,
781        );
782        Ok(())
783    }
784
785    fn unsubscribe_book_deltas(&mut self, cmd: &UnsubscribeBookDeltas) -> anyhow::Result<()> {
786        let instrument_id = cmd.instrument_id;
787
788        self.book_instruments.remove(&instrument_id);
789
790        let ws = self.ws.clone();
791        self.spawn_ws(
792            async move {
793                ws.unsubscribe_book(instrument_id)
794                    .await
795                    .map_err(|e| anyhow::anyhow!("{e}"))
796            },
797            "unsubscribe book",
798        );
799
800        Ok(())
801    }
802
803    fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
804        let instrument_id = cmd.instrument_id;
805        let ws = self.ws.clone();
806
807        self.quote_instruments.remove(&instrument_id);
808
809        self.spawn_ws(
810            async move {
811                ws.unsubscribe_quotes(instrument_id)
812                    .await
813                    .map_err(|e| anyhow::anyhow!("{e}"))
814            },
815            "unsubscribe quotes",
816        );
817
818        Ok(())
819    }
820
821    fn unsubscribe_trades(&mut self, cmd: &UnsubscribeTrades) -> anyhow::Result<()> {
822        let instrument_id = cmd.instrument_id;
823        let ws = self.ws.clone();
824
825        self.spawn_ws(
826            async move {
827                ws.unsubscribe_trades(instrument_id)
828                    .await
829                    .map_err(|e| anyhow::anyhow!("{e}"))
830            },
831            "unsubscribe trades",
832        );
833
834        Ok(())
835    }
836
837    fn unsubscribe_mark_prices(&mut self, cmd: &UnsubscribeMarkPrices) -> anyhow::Result<()> {
838        let instrument_id = cmd.instrument_id;
839        let ws = self.ws.clone();
840
841        self.spawn_ws(
842            async move {
843                ws.unsubscribe_mark_price(instrument_id)
844                    .await
845                    .map_err(|e| anyhow::anyhow!("{e}"))
846            },
847            "unsubscribe mark price",
848        );
849
850        Ok(())
851    }
852
853    fn unsubscribe_index_prices(&mut self, cmd: &UnsubscribeIndexPrices) -> anyhow::Result<()> {
854        let instrument_id = cmd.instrument_id;
855        let ws = self.ws.clone();
856
857        self.spawn_ws(
858            async move {
859                ws.unsubscribe_index_price(instrument_id)
860                    .await
861                    .map_err(|e| anyhow::anyhow!("{e}"))
862            },
863            "unsubscribe index price",
864        );
865
866        Ok(())
867    }
868
869    fn unsubscribe_funding_rates(&mut self, cmd: &UnsubscribeFundingRates) -> anyhow::Result<()> {
870        let instrument_id = cmd.instrument_id;
871        let ws = self.ws.clone();
872
873        self.spawn_ws(
874            async move {
875                ws.unsubscribe_funding_rate(instrument_id)
876                    .await
877                    .map_err(|e| anyhow::anyhow!("{e}"))
878            },
879            "unsubscribe funding rate",
880        );
881
882        Ok(())
883    }
884
885    fn unsubscribe_bars(&mut self, _cmd: &UnsubscribeBars) -> anyhow::Result<()> {
886        Ok(())
887    }
888
889    fn unsubscribe_instrument_status(
890        &mut self,
891        _cmd: &UnsubscribeInstrumentStatus,
892    ) -> anyhow::Result<()> {
893        Ok(())
894    }
895
896    fn request_instruments(&self, request: RequestInstruments) -> anyhow::Result<()> {
897        let http = self.http.clone();
898        let sender = self.data_sender.clone();
899        let instruments_cache = self.instruments.clone();
900        let request_id = request.request_id;
901        let client_id = request.client_id.unwrap_or(self.client_id);
902        let venue = *KRAKEN_VENUE;
903        let start_nanos = datetime_to_unix_nanos(request.start);
904        let end_nanos = datetime_to_unix_nanos(request.end);
905        let params = request.params;
906        let clock = self.clock;
907
908        self.spawn_command(async move {
909            match http.request_instruments().await {
910                Ok(instruments) => {
911                    instruments_cache.rcu(|m| {
912                        for instrument in &instruments {
913                            m.insert(instrument.id(), instrument.clone());
914                        }
915                    });
916                    http.cache_instruments(&instruments);
917
918                    let response = DataResponse::Instruments(InstrumentsResponse::new(
919                        request_id,
920                        client_id,
921                        venue,
922                        instruments,
923                        start_nanos,
924                        end_nanos,
925                        clock.get_time_ns(),
926                        params,
927                    ));
928
929                    if let Err(e) = sender.send(DataEvent::Response(response)) {
930                        log::error!("Failed to send instruments response: {e}");
931                    }
932                }
933                Err(e) => log::error!("Instruments request failed: {e:?}"),
934            }
935        });
936
937        Ok(())
938    }
939
940    fn request_instrument(&self, request: RequestInstrument) -> anyhow::Result<()> {
941        let http = self.http.clone();
942        let sender = self.data_sender.clone();
943        let instruments = self.instruments.clone();
944        let instrument_id = request.instrument_id;
945        let request_id = request.request_id;
946        let client_id = request.client_id.unwrap_or(self.client_id);
947        let start_nanos = datetime_to_unix_nanos(request.start);
948        let end_nanos = datetime_to_unix_nanos(request.end);
949        let params = request.params;
950        let clock = self.clock;
951
952        self.spawn_command(async move {
953            match http.request_instruments().await {
954                Ok(all_instruments) => {
955                    instruments.rcu(|m| {
956                        for instrument in &all_instruments {
957                            m.insert(instrument.id(), instrument.clone());
958                        }
959                    });
960                    http.cache_instruments(&all_instruments);
961
962                    let instrument = all_instruments
963                        .into_iter()
964                        .find(|i| i.id() == instrument_id);
965
966                    if let Some(instrument) = instrument {
967                        let response = DataResponse::Instrument(Box::new(InstrumentResponse::new(
968                            request_id,
969                            client_id,
970                            instrument.id(),
971                            instrument,
972                            start_nanos,
973                            end_nanos,
974                            clock.get_time_ns(),
975                            params,
976                        )));
977
978                        if let Err(e) = sender.send(DataEvent::Response(response)) {
979                            log::error!("Failed to send instrument response: {e}");
980                        }
981                    } else {
982                        log::error!("Instrument not found: {instrument_id}");
983                    }
984                }
985                Err(e) => log::error!("Instrument request failed: {e:?}"),
986            }
987        });
988
989        Ok(())
990    }
991
992    fn request_trades(&self, request: RequestTrades) -> anyhow::Result<()> {
993        let http = self.http.clone();
994        let sender = self.data_sender.clone();
995        let instrument_id = request.instrument_id;
996        let start = request.start;
997        let end = request.end;
998        let limit = request.limit.map(|n| n.get() as u64);
999        let request_id = request.request_id;
1000        let client_id = request.client_id.unwrap_or(self.client_id);
1001        let params = request.params;
1002        let clock = self.clock;
1003        let start_nanos = datetime_to_unix_nanos(start);
1004        let end_nanos = datetime_to_unix_nanos(end);
1005
1006        self.spawn_command(async move {
1007            match http.request_trades(instrument_id, start, end, limit).await {
1008                Ok(trades) => {
1009                    let response = DataResponse::Trades(TradesResponse::new(
1010                        request_id,
1011                        client_id,
1012                        instrument_id,
1013                        trades,
1014                        start_nanos,
1015                        end_nanos,
1016                        clock.get_time_ns(),
1017                        params,
1018                    ));
1019
1020                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1021                        log::error!("Failed to send trades response: {e}");
1022                    }
1023                }
1024                Err(e) => log::error!("Trades request failed: {e:?}"),
1025            }
1026        });
1027
1028        Ok(())
1029    }
1030
1031    fn request_bars(&self, request: RequestBars) -> anyhow::Result<()> {
1032        let http = self.http.clone();
1033        let sender = self.data_sender.clone();
1034        let bar_type = request.bar_type;
1035        let start = request.start;
1036        let end = request.end;
1037        let limit = request.limit.map(|n| n.get() as u64);
1038        let request_id = request.request_id;
1039        let client_id = request.client_id.unwrap_or(self.client_id);
1040        let params = request.params;
1041        let clock = self.clock;
1042        let start_nanos = datetime_to_unix_nanos(start);
1043        let end_nanos = datetime_to_unix_nanos(end);
1044
1045        self.spawn_command(async move {
1046            match http.request_bars(bar_type, start, end, limit).await {
1047                Ok(bars) => {
1048                    let response = DataResponse::Bars(BarsResponse::new(
1049                        request_id,
1050                        client_id,
1051                        bar_type,
1052                        bars,
1053                        start_nanos,
1054                        end_nanos,
1055                        clock.get_time_ns(),
1056                        params,
1057                    ));
1058
1059                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1060                        log::error!("Failed to send bars response: {e}");
1061                    }
1062                }
1063                Err(e) => log::error!("Bars request failed: {e:?}"),
1064            }
1065        });
1066
1067        Ok(())
1068    }
1069
1070    fn request_book_snapshot(&self, request: RequestBookSnapshot) -> anyhow::Result<()> {
1071        let http = self.http.clone();
1072        let sender = self.data_sender.clone();
1073        let instrument_id = request.instrument_id;
1074        let depth = request.depth.map(|n| n.get() as u32);
1075        let request_id = request.request_id;
1076        let client_id = request.client_id.unwrap_or(self.client_id);
1077        let params = request.params;
1078        let clock = self.clock;
1079
1080        self.spawn_command(async move {
1081            match http.request_book_snapshot(instrument_id, depth).await {
1082                Ok(book) => {
1083                    let response = DataResponse::Book(BookResponse::new(
1084                        request_id,
1085                        client_id,
1086                        instrument_id,
1087                        book,
1088                        None,
1089                        None,
1090                        clock.get_time_ns(),
1091                        params,
1092                    ));
1093
1094                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1095                        log::error!("Failed to send book snapshot response: {e}");
1096                    }
1097                }
1098                Err(e) => log::error!("Book snapshot request failed: {e:?}"),
1099            }
1100        });
1101
1102        Ok(())
1103    }
1104
1105    fn request_funding_rates(&self, request: RequestFundingRates) -> anyhow::Result<()> {
1106        let http = self.http.clone();
1107        let sender = self.data_sender.clone();
1108        let instrument_id = request.instrument_id;
1109        let start = request.start;
1110        let end = request.end;
1111        let limit = request.limit.map(|n| n.get());
1112        let request_id = request.request_id;
1113        let client_id = request.client_id.unwrap_or(self.client_id);
1114        let start_nanos = datetime_to_unix_nanos(start);
1115        let end_nanos = datetime_to_unix_nanos(end);
1116        let params = request.params;
1117        let clock = self.clock;
1118
1119        self.spawn_command(async move {
1120            match http
1121                .request_funding_rates(instrument_id, start, end, limit)
1122                .await
1123            {
1124                Ok(rates) => {
1125                    let response = DataResponse::FundingRates(FundingRatesResponse::new(
1126                        request_id,
1127                        client_id,
1128                        instrument_id,
1129                        rates,
1130                        start_nanos,
1131                        end_nanos,
1132                        clock.get_time_ns(),
1133                        params,
1134                    ));
1135
1136                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1137                        log::error!("Failed to send funding rates response: {e}");
1138                    }
1139                }
1140                Err(e) => log::error!("Funding rates request failed: {e:?}"),
1141            }
1142        });
1143
1144        Ok(())
1145    }
1146}
1147
1148#[cfg(test)]
1149mod tests {
1150    use nautilus_common::{live::runner::set_data_event_sender, messages::DataEvent};
1151    use rstest::rstest;
1152
1153    use super::*;
1154    use crate::{
1155        common::{consts::KRAKEN_CLIENT_ID, enums::KrakenProductType},
1156        config::KrakenDataClientConfig,
1157    };
1158
1159    fn setup_test_env() {
1160        let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1161        set_data_event_sender(sender);
1162    }
1163
1164    #[rstest]
1165    fn test_futures_data_client_new() {
1166        setup_test_env();
1167        let config = KrakenDataClientConfig {
1168            product_type: KrakenProductType::Futures,
1169            ..Default::default()
1170        };
1171        let client = KrakenFuturesDataClient::new(*KRAKEN_CLIENT_ID, config);
1172        assert!(client.is_ok());
1173
1174        let client = client.unwrap();
1175        assert_eq!(client.client_id(), *KRAKEN_CLIENT_ID);
1176        assert_eq!(client.venue(), Some(*KRAKEN_VENUE));
1177        assert!(!client.is_connected());
1178        assert!(client.is_disconnected());
1179        assert!(client.instruments().is_empty());
1180    }
1181
1182    #[rstest]
1183    fn test_futures_data_client_start_stop() {
1184        setup_test_env();
1185        let config = KrakenDataClientConfig {
1186            product_type: KrakenProductType::Futures,
1187            ..Default::default()
1188        };
1189        let mut client = KrakenFuturesDataClient::new(*KRAKEN_CLIENT_ID, config).unwrap();
1190
1191        assert!(client.start().is_ok());
1192        assert!(client.stop().is_ok());
1193        assert!(client.is_disconnected());
1194    }
1195}