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