1use 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#[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 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 #[must_use]
159 pub fn instruments(&self) -> Vec<InstrumentAny> {
160 self.instruments.load().values().cloned().collect()
161 }
162
163 #[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 "e_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}