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