Skip to main content

nautilus_kraken/websocket/futures/
handler.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//! WebSocket message handler for Kraken Futures.
17
18use std::{
19    collections::VecDeque,
20    sync::{
21        Arc,
22        atomic::{AtomicBool, Ordering},
23    },
24};
25
26use nautilus_core::string::secret::SecretString;
27use nautilus_network::{
28    RECONNECTED,
29    websocket::{SubscriptionState, WebSocketClient},
30};
31use serde::Deserialize;
32use serde_json::Value;
33use tokio_tungstenite::tungstenite::Message;
34use ustr::Ustr;
35
36use super::messages::{
37    KrakenFuturesBookDelta, KrakenFuturesBookSnapshot, KrakenFuturesChannel,
38    KrakenFuturesFillsDelta, KrakenFuturesMessageType, KrakenFuturesOpenOrdersCancel,
39    KrakenFuturesOpenOrdersDelta, KrakenFuturesTickerData, KrakenFuturesTradeData,
40    KrakenFuturesWsMessage, classify_futures_message,
41};
42use crate::common::consts::KRAKEN_RATE_LIMIT_KEY_SUBSCRIPTION;
43
44/// Commands sent from the outer client to the inner message handler.
45#[derive(Debug)]
46pub enum FuturesHandlerCommand {
47    SetClient(WebSocketClient),
48    Disconnect,
49    Subscribe { payload: SecretString },
50    Unsubscribe { payload: SecretString },
51    RequestChallenge { payload: SecretString },
52}
53
54/// WebSocket message handler for Kraken Futures.
55pub struct FuturesFeedHandler {
56    signal: Arc<AtomicBool>,
57    inner: Option<WebSocketClient>,
58    cmd_rx: tokio::sync::mpsc::UnboundedReceiver<FuturesHandlerCommand>,
59    raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
60    subscriptions: SubscriptionState,
61    pending_messages: VecDeque<KrakenFuturesWsMessage>,
62}
63
64impl FuturesFeedHandler {
65    /// Creates a new [`FuturesFeedHandler`] instance.
66    pub fn new(
67        signal: Arc<AtomicBool>,
68        cmd_rx: tokio::sync::mpsc::UnboundedReceiver<FuturesHandlerCommand>,
69        raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
70        subscriptions: SubscriptionState,
71    ) -> Self {
72        Self {
73            signal,
74            inner: None,
75            cmd_rx,
76            raw_rx,
77            subscriptions,
78            pending_messages: VecDeque::new(),
79        }
80    }
81
82    pub fn is_stopped(&self) -> bool {
83        self.signal.load(Ordering::Relaxed)
84    }
85
86    fn is_subscribed(&self, channel: KrakenFuturesChannel, symbol: &Ustr) -> bool {
87        let channel_ustr = Ustr::from(channel.as_ref());
88        self.subscriptions.is_subscribed(&channel_ustr, symbol)
89    }
90
91    /// Processes messages and commands, returning when stopped or stream ends.
92    pub async fn next(&mut self) -> Option<KrakenFuturesWsMessage> {
93        if let Some(msg) = self.pending_messages.pop_front() {
94            return Some(msg);
95        }
96
97        loop {
98            tokio::select! {
99                Some(cmd) = self.cmd_rx.recv() => {
100                    match cmd {
101                        FuturesHandlerCommand::SetClient(client) => {
102                            log::debug!("WebSocketClient received by futures handler");
103                            self.inner = Some(client);
104                        }
105                        FuturesHandlerCommand::Disconnect => {
106                            log::debug!("Disconnect command received");
107
108                            if let Some(client) = self.inner.take() {
109                                client.disconnect().await;
110                            }
111                            return None;
112                        }
113                        FuturesHandlerCommand::Subscribe { payload }
114                        | FuturesHandlerCommand::Unsubscribe { payload } => {
115                            if let Some(ref client) = self.inner
116                                && let Err(e) = client.send_text(payload.expose_secret().to_owned(), Some(KRAKEN_RATE_LIMIT_KEY_SUBSCRIPTION.as_slice())).await
117                            {
118                                log::error!("Failed to send text: {e}");
119                            }
120                        }
121                        FuturesHandlerCommand::RequestChallenge { payload } => {
122                            if let Some(ref client) = self.inner
123                                && let Err(e) = client.send_text(payload.expose_secret().to_owned(), None).await
124                            {
125                                log::error!("Failed to send text: {e}");
126                            }
127                        }
128                    }
129                }
130
131                msg = self.raw_rx.recv() => {
132                    let msg = match msg {
133                        Some(msg) => msg,
134                        None => {
135                            log::debug!("WebSocket stream closed");
136                            return None;
137                        }
138                    };
139
140                    if self.signal.load(Ordering::Relaxed) {
141                        log::debug!("Stop signal received");
142                        return None;
143                    }
144
145                    match &msg {
146                        Message::Ping(data) => {
147                            let len = data.len();
148                            log::trace!("Received ping frame with {len} bytes");
149
150                            if let Some(client) = &self.inner
151                                && let Err(e) = client.send_pong(data.to_vec()).await
152                            {
153                                log::warn!("Failed to send pong frame: {e}");
154                            }
155                            continue;
156                        }
157                        Message::Pong(_) => {
158                            log::debug!("Received pong from server");
159                            continue;
160                        }
161                        Message::Close(_) => {
162                            log::debug!("WebSocket connection closed");
163                            return None;
164                        }
165                        Message::Frame(_) => {
166                            log::trace!("Received raw frame");
167                            continue;
168                        }
169                        _ => {}
170                    }
171
172                    let text: &str = match &msg {
173                        Message::Text(text) => text,
174                        Message::Binary(data) => match std::str::from_utf8(data) {
175                            Ok(s) => s,
176                            Err(_) => continue,
177                        },
178                        _ => continue,
179                    };
180
181                    if text == RECONNECTED {
182                        log::debug!("Received WebSocket reconnected signal");
183                        return Some(KrakenFuturesWsMessage::Reconnected);
184                    }
185
186                    self.parse_message(text);
187
188                    if let Some(msg) = self.pending_messages.pop_front() {
189                        return Some(msg);
190                    }
191                }
192            }
193        }
194    }
195
196    fn parse_message(&mut self, text: &str) {
197        let value: Value = match serde_json::from_str(text) {
198            Ok(v) => v,
199            Err(e) => {
200                log::debug!("Failed to parse message as JSON: {e}");
201                return;
202            }
203        };
204
205        match classify_futures_message(&value) {
206            KrakenFuturesMessageType::OpenOrdersSnapshot => {
207                log::debug!(
208                    "Skipping open_orders_snapshot (REST reconciliation handles initial state)"
209                );
210            }
211            KrakenFuturesMessageType::OpenOrdersCancel => {
212                self.handle_open_orders_cancel_text(text);
213            }
214            KrakenFuturesMessageType::OpenOrdersDelta => {
215                self.handle_open_orders_delta_text(text);
216            }
217            KrakenFuturesMessageType::FillsSnapshot => {
218                log::debug!("Skipping fills_snapshot (REST reconciliation handles initial state)");
219            }
220            KrakenFuturesMessageType::FillsDelta => {
221                self.handle_fills_delta_text(text);
222            }
223            KrakenFuturesMessageType::Ticker => {
224                self.handle_ticker_message_text(text);
225            }
226            KrakenFuturesMessageType::TradeSnapshot => {
227                log::debug!("Skipping trade_snapshot (only streaming live trades)");
228            }
229            KrakenFuturesMessageType::Trade => {
230                self.handle_trade_message_text(text);
231            }
232            KrakenFuturesMessageType::BookSnapshot => {
233                self.handle_book_snapshot_text(text);
234            }
235            KrakenFuturesMessageType::BookDelta => {
236                self.handle_book_delta_text(text);
237            }
238            KrakenFuturesMessageType::Info => {
239                log::debug!("Received info message: {text}");
240            }
241            KrakenFuturesMessageType::Pong => {
242                log::debug!("Received text pong response");
243            }
244            KrakenFuturesMessageType::Subscribed => {
245                log::debug!("Subscription confirmed: {text}");
246            }
247            KrakenFuturesMessageType::Unsubscribed => {
248                log::debug!("Unsubscription confirmed: {text}");
249            }
250            KrakenFuturesMessageType::Challenge => {
251                self.handle_challenge_response_text(text);
252            }
253            KrakenFuturesMessageType::Heartbeat => {
254                log::trace!("Heartbeat received");
255            }
256            KrakenFuturesMessageType::Error => {
257                let message = value
258                    .get("message")
259                    .and_then(|v| v.as_str())
260                    .unwrap_or("Unknown error");
261                log::warn!("Kraken Futures WebSocket error: {message}");
262            }
263            KrakenFuturesMessageType::Alert => {
264                let message = value
265                    .get("message")
266                    .and_then(|v| v.as_str())
267                    .unwrap_or("Unknown alert");
268                log::warn!("Kraken Futures WebSocket alert: {message}");
269            }
270            KrakenFuturesMessageType::Unknown => {
271                log::warn!("Unhandled futures message: {text}");
272            }
273        }
274    }
275
276    fn handle_challenge_response_text(&mut self, text: &str) {
277        #[derive(Deserialize)]
278        struct ChallengeResponse {
279            message: String,
280        }
281
282        match serde_json::from_str::<ChallengeResponse>(text) {
283            Ok(response) => {
284                let len = response.message.len();
285                log::debug!("Challenge received, length: {len}");
286
287                self.pending_messages
288                    .push_back(KrakenFuturesWsMessage::Challenge(response.message));
289            }
290            Err(e) => {
291                log::error!("Failed to parse challenge response: {e}");
292            }
293        }
294    }
295
296    fn handle_ticker_message_text(&mut self, text: &str) {
297        let ticker = match serde_json::from_str::<KrakenFuturesTickerData>(text) {
298            Ok(t) => t,
299            Err(e) => {
300                log::debug!("Failed to parse ticker: {e}");
301                return;
302            }
303        };
304
305        self.pending_messages
306            .push_back(KrakenFuturesWsMessage::Ticker(ticker));
307    }
308
309    fn handle_trade_message_text(&mut self, text: &str) {
310        let trade = match serde_json::from_str::<KrakenFuturesTradeData>(text) {
311            Ok(t) => t,
312            Err(e) => {
313                log::warn!("Failed to parse trade: {e}");
314                return;
315            }
316        };
317
318        if !self.is_subscribed(KrakenFuturesChannel::Trades, &trade.product_id) {
319            log::debug!(
320                "Received trade for unsubscribed product: {}",
321                trade.product_id
322            );
323            return;
324        }
325
326        self.pending_messages
327            .push_back(KrakenFuturesWsMessage::Trade(trade));
328    }
329
330    fn handle_book_snapshot_text(&mut self, text: &str) {
331        let snapshot = match serde_json::from_str::<KrakenFuturesBookSnapshot>(text) {
332            Ok(s) => s,
333            Err(e) => {
334                log::warn!("Failed to parse book snapshot: {e}");
335                return;
336            }
337        };
338
339        let has_book = self.is_subscribed(KrakenFuturesChannel::Book, &snapshot.product_id);
340        let has_quotes = self.is_subscribed(KrakenFuturesChannel::Quotes, &snapshot.product_id);
341
342        if !has_book && !has_quotes {
343            log::debug!(
344                "Received book snapshot for unsubscribed product: {}",
345                snapshot.product_id
346            );
347            return;
348        }
349
350        self.pending_messages
351            .push_back(KrakenFuturesWsMessage::BookSnapshot(snapshot));
352    }
353
354    fn handle_book_delta_text(&mut self, text: &str) {
355        let delta = match serde_json::from_str::<KrakenFuturesBookDelta>(text) {
356            Ok(d) => d,
357            Err(e) => {
358                log::warn!("Failed to parse book delta: {e}");
359                return;
360            }
361        };
362
363        let has_book = self.is_subscribed(KrakenFuturesChannel::Book, &delta.product_id);
364        let has_quotes = self.is_subscribed(KrakenFuturesChannel::Quotes, &delta.product_id);
365
366        if !has_book && !has_quotes {
367            log::debug!(
368                "Received book delta for unsubscribed product: {}",
369                delta.product_id
370            );
371            return;
372        }
373
374        self.pending_messages
375            .push_back(KrakenFuturesWsMessage::BookDelta(delta));
376    }
377
378    fn handle_open_orders_delta_text(&mut self, text: &str) {
379        let delta = match serde_json::from_str::<KrakenFuturesOpenOrdersDelta>(text) {
380            Ok(d) => d,
381            Err(e) => {
382                log::error!("Failed to parse open_orders delta: {e}");
383                return;
384            }
385        };
386
387        log::debug!(
388            "Received open_orders delta: order_id={}, is_cancel={}, reason={:?}",
389            delta.order.order_id,
390            delta.is_cancel,
391            delta.reason
392        );
393
394        self.pending_messages
395            .push_back(KrakenFuturesWsMessage::OpenOrdersDelta(delta));
396    }
397
398    fn handle_open_orders_cancel_text(&mut self, text: &str) {
399        let cancel = match serde_json::from_str::<KrakenFuturesOpenOrdersCancel>(text) {
400            Ok(c) => c,
401            Err(e) => {
402                log::error!("Failed to parse open_orders cancel: {e}");
403                return;
404            }
405        };
406
407        log::debug!(
408            "Received open_orders cancel: order_id={}, cli_ord_id={:?}, reason={:?}",
409            cancel.order_id,
410            cancel.cli_ord_id,
411            cancel.reason
412        );
413
414        self.pending_messages
415            .push_back(KrakenFuturesWsMessage::OpenOrdersCancel(cancel));
416    }
417
418    fn handle_fills_delta_text(&mut self, text: &str) {
419        let delta = match serde_json::from_str::<KrakenFuturesFillsDelta>(text) {
420            Ok(d) => d,
421            Err(e) => {
422                log::error!("Failed to parse fills delta: {e}");
423                return;
424            }
425        };
426
427        log::debug!("Received fills delta: fill_count={}", delta.fills.len());
428
429        self.pending_messages
430            .push_back(KrakenFuturesWsMessage::FillsDelta(delta));
431    }
432}
433
434#[cfg(test)]
435mod tests {
436    use nautilus_core::string::secret::REDACTED;
437    use rstest::rstest;
438    use rust_decimal_macros::dec;
439
440    use super::*;
441
442    fn create_test_handler() -> FuturesFeedHandler {
443        let signal = Arc::new(AtomicBool::new(false));
444        let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
445        let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
446        let subscriptions = SubscriptionState::new(':');
447
448        FuturesFeedHandler::new(signal, cmd_rx, raw_rx, subscriptions)
449    }
450
451    #[rstest]
452    fn test_command_debug_redacts_payload() {
453        let command = FuturesHandlerCommand::RequestChallenge {
454            payload: SecretString::from("private-challenge-payload"),
455        };
456
457        assert_eq!(
458            format!("{command:?}"),
459            format!("RequestChallenge {{ payload: {REDACTED} }}")
460        );
461    }
462
463    #[rstest]
464    fn test_parse_ticker_emits_ticker_message() {
465        let mut handler = create_test_handler();
466        let json = include_str!("../../../test_data/ws_futures_ticker.json");
467
468        handler.parse_message(json);
469
470        assert_eq!(handler.pending_messages.len(), 1);
471        let msg = handler.pending_messages.pop_front().unwrap();
472        let KrakenFuturesWsMessage::Ticker(ticker) = msg else {
473            panic!("Expected Ticker message, was {msg:?}");
474        };
475        assert_eq!(ticker.product_id, Ustr::from("PI_XBTUSD"));
476        assert_eq!(ticker.bid, Some(dec!(21978.5)));
477        assert_eq!(ticker.ask, Some(dec!(21987)));
478    }
479
480    #[rstest]
481    fn test_parse_trade_emits_trade_message() {
482        let mut handler = create_test_handler();
483        handler.subscriptions.mark_subscribe("trades:PI_XBTUSD");
484        handler.subscriptions.confirm_subscribe("trades:PI_XBTUSD");
485
486        let json = include_str!("../../../test_data/ws_futures_trade.json");
487
488        handler.parse_message(json);
489
490        assert_eq!(handler.pending_messages.len(), 1);
491        let msg = handler.pending_messages.pop_front().unwrap();
492        let KrakenFuturesWsMessage::Trade(trade) = msg else {
493            panic!("Expected Trade message, was {msg:?}");
494        };
495        assert_eq!(trade.product_id, Ustr::from("PI_XBTUSD"));
496        assert_eq!(trade.price, dec!(34969.5));
497        assert_eq!(trade.qty, dec!(15000));
498    }
499
500    #[rstest]
501    fn test_parse_trade_filters_unsubscribed() {
502        let mut handler = create_test_handler();
503        let json = include_str!("../../../test_data/ws_futures_trade.json");
504
505        handler.parse_message(json);
506
507        assert!(
508            handler.pending_messages.is_empty(),
509            "Trade for unsubscribed product should be filtered"
510        );
511    }
512
513    #[rstest]
514    fn test_parse_book_snapshot_emits_book_snapshot() {
515        let mut handler = create_test_handler();
516        handler.subscriptions.mark_subscribe("book:PI_XBTUSD");
517        handler.subscriptions.confirm_subscribe("book:PI_XBTUSD");
518
519        let json = include_str!("../../../test_data/ws_futures_book_snapshot.json");
520
521        handler.parse_message(json);
522
523        assert_eq!(handler.pending_messages.len(), 1);
524        let msg = handler.pending_messages.pop_front().unwrap();
525        let KrakenFuturesWsMessage::BookSnapshot(snapshot) = msg else {
526            panic!("Expected BookSnapshot message, was {msg:?}");
527        };
528        assert_eq!(snapshot.product_id, Ustr::from("PI_XBTUSD"));
529        assert_eq!(snapshot.bids.len(), 2);
530        assert_eq!(snapshot.asks.len(), 2);
531    }
532
533    #[rstest]
534    fn test_parse_book_snapshot_filters_unsubscribed() {
535        let mut handler = create_test_handler();
536        let json = include_str!("../../../test_data/ws_futures_book_snapshot.json");
537
538        handler.parse_message(json);
539
540        assert!(
541            handler.pending_messages.is_empty(),
542            "Book snapshot for unsubscribed product should be filtered"
543        );
544    }
545
546    #[rstest]
547    fn test_parse_book_delta_emits_book_delta() {
548        let mut handler = create_test_handler();
549        handler.subscriptions.mark_subscribe("book:PI_XBTUSD");
550        handler.subscriptions.confirm_subscribe("book:PI_XBTUSD");
551
552        let json = include_str!("../../../test_data/ws_futures_book_delta.json");
553
554        handler.parse_message(json);
555
556        assert_eq!(handler.pending_messages.len(), 1);
557        let msg = handler.pending_messages.pop_front().unwrap();
558        let KrakenFuturesWsMessage::BookDelta(delta) = msg else {
559            panic!("Expected BookDelta message, was {msg:?}");
560        };
561        assert_eq!(delta.product_id, Ustr::from("PI_XBTUSD"));
562        assert_eq!(delta.price, dec!(34981));
563    }
564
565    #[rstest]
566    fn test_parse_book_delta_preserves_decimal_precision() {
567        let mut handler = create_test_handler();
568        handler.subscriptions.mark_subscribe("book:PI_XBTUSD");
569        handler.subscriptions.confirm_subscribe("book:PI_XBTUSD");
570        let json = include_str!("../../../test_data/ws_futures_book_precision.json");
571
572        handler.parse_message(json);
573
574        let message = handler.pending_messages.pop_front().unwrap();
575        let KrakenFuturesWsMessage::BookDelta(delta) = message else {
576            panic!("Expected BookDelta message, was {message:?}");
577        };
578        assert_eq!(delta.price, dec!(123456789.123456789));
579        assert_eq!(delta.qty, dec!(0.1234567890123456789012345678));
580    }
581
582    #[rstest]
583    fn test_parse_book_delta_filters_unsubscribed() {
584        let mut handler = create_test_handler();
585        let json = include_str!("../../../test_data/ws_futures_book_delta.json");
586
587        handler.parse_message(json);
588
589        assert!(
590            handler.pending_messages.is_empty(),
591            "Book delta for unsubscribed product should be filtered"
592        );
593    }
594
595    #[rstest]
596    fn test_parse_open_orders_cancel_emits_cancel() {
597        let mut handler = create_test_handler();
598        let json = include_str!("../../../test_data/ws_futures_open_orders_cancel.json");
599
600        handler.parse_message(json);
601
602        assert_eq!(handler.pending_messages.len(), 1);
603        let msg = handler.pending_messages.pop_front().unwrap();
604        let KrakenFuturesWsMessage::OpenOrdersCancel(cancel) = msg else {
605            panic!("Expected OpenOrdersCancel message, was {msg:?}");
606        };
607        assert_eq!(cancel.order_id, "660c6b23-8007-48c1-a7c9-4893f4572e8c");
608        assert!(cancel.is_cancel);
609    }
610
611    #[rstest]
612    fn test_parse_open_orders_delta_emits_delta() {
613        let mut handler = create_test_handler();
614        let json = include_str!("../../../test_data/ws_futures_open_orders_delta.json");
615
616        handler.parse_message(json);
617
618        assert_eq!(handler.pending_messages.len(), 1);
619        let msg = handler.pending_messages.pop_front().unwrap();
620        let KrakenFuturesWsMessage::OpenOrdersDelta(delta) = msg else {
621            panic!("Expected OpenOrdersDelta message, was {msg:?}");
622        };
623        assert_eq!(delta.order.instrument, Ustr::from("PI_XBTUSD"));
624        assert!(!delta.is_cancel);
625    }
626
627    #[rstest]
628    fn test_parse_fills_delta_emits_fills() {
629        let mut handler = create_test_handler();
630        let json = include_str!("../../../test_data/ws_futures_fills_delta.json");
631
632        handler.parse_message(json);
633
634        assert_eq!(handler.pending_messages.len(), 1);
635        let msg = handler.pending_messages.pop_front().unwrap();
636        let KrakenFuturesWsMessage::FillsDelta(fills) = msg else {
637            panic!("Expected FillsDelta message, was {msg:?}");
638        };
639        assert_eq!(fills.fills.len(), 1);
640        assert_eq!(
641            fills.fills[0].fill_id,
642            "6a22a3fb-e18e-4e76-b841-8689735c9158"
643        );
644    }
645
646    #[rstest]
647    fn test_parse_challenge_emits_challenge_message() {
648        let mut handler = create_test_handler();
649        let json = r#"{"event":"challenge","message":"server-challenge-abc"}"#;
650
651        handler.parse_message(json);
652
653        assert_eq!(handler.pending_messages.len(), 1);
654        let msg = handler.pending_messages.pop_front().unwrap();
655        let KrakenFuturesWsMessage::Challenge(challenge) = msg else {
656            panic!("Expected Challenge message, was {msg:?}");
657        };
658        assert_eq!(challenge, "server-challenge-abc");
659    }
660
661    #[rstest]
662    fn test_heartbeat_produces_no_message() {
663        let mut handler = create_test_handler();
664        let json = r#"{"feed":"heartbeat","time":1700000000000}"#;
665
666        handler.parse_message(json);
667
668        assert!(handler.pending_messages.is_empty());
669    }
670
671    #[rstest]
672    fn test_info_event_produces_no_message() {
673        let mut handler = create_test_handler();
674        let json = r#"{"event":"info","version":1}"#;
675
676        handler.parse_message(json);
677
678        assert!(handler.pending_messages.is_empty());
679    }
680}