nautilus_kraken/websocket/futures/
handler.rs1use 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#[derive(Debug)]
45pub enum FuturesHandlerCommand {
46 SetClient(WebSocketClient),
47 Disconnect,
48 Subscribe { payload: String },
49 Unsubscribe { payload: String },
50 RequestChallenge { payload: String },
51}
52
53pub 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 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 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}