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,
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: tokio::sync::mpsc::UnboundedSender<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
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 #[must_use]
155 pub fn instruments(&self) -> Vec<InstrumentAny> {
156 self.instruments.load().values().cloned().collect()
157 }
158
159 #[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 "e_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}