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