1use std::{
19 collections::VecDeque,
20 sync::{
21 Arc,
22 atomic::{AtomicBool, Ordering},
23 },
24};
25
26#[cfg(test)]
27use nautilus_core::string::secret::REDACTED;
28use nautilus_core::string::secret::SecretString;
29use nautilus_network::{
30 RECONNECTED,
31 websocket::{SubscriptionState, WebSocketClient},
32};
33use serde::Deserialize;
34use serde_json::{Value, value::RawValue};
35use tokio_tungstenite::tungstenite::Message;
36
37use super::{
38 enums::{KrakenWsChannel, KrakenWsMessageType},
39 messages::{
40 KrakenSpotWsMessage, KrakenWsBookData, KrakenWsExecutionData, KrakenWsOhlcData,
41 KrakenWsRawMessage, KrakenWsResponse, KrakenWsTickerData, KrakenWsTradeData,
42 },
43 parse::parse_order_response,
44};
45use crate::{
46 common::consts::{KRAKEN_RATE_LIMIT_KEY_ORDER, KRAKEN_RATE_LIMIT_KEY_SUBSCRIPTION},
47 websocket::spot_v2::level_3::messages::{KrakenL3Snapshot, KrakenL3UpdateData},
48};
49
50#[derive(Debug)]
52pub enum SpotHandlerCommand {
53 SetClient(WebSocketClient),
54 Disconnect,
55 Subscribe { payload: SecretString },
56 Unsubscribe { payload: SecretString },
57 Ping { payload: SecretString },
58 SendOrderRequest { req_id: u64, payload: SecretString },
59}
60
61pub(super) struct SpotFeedHandler {
63 signal: Arc<AtomicBool>,
64 inner: Option<WebSocketClient>,
65 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<SpotHandlerCommand>,
66 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
67 subscriptions: SubscriptionState,
68 pending_messages: VecDeque<KrakenSpotWsMessage>,
69}
70
71impl SpotFeedHandler {
72 pub(super) fn new(
74 signal: Arc<AtomicBool>,
75 cmd_rx: tokio::sync::mpsc::UnboundedReceiver<SpotHandlerCommand>,
76 raw_rx: tokio::sync::mpsc::UnboundedReceiver<Message>,
77 subscriptions: SubscriptionState,
78 ) -> Self {
79 Self {
80 signal,
81 inner: None,
82 cmd_rx,
83 raw_rx,
84 subscriptions,
85 pending_messages: VecDeque::new(),
86 }
87 }
88
89 pub(super) fn is_stopped(&self) -> bool {
90 self.signal.load(Ordering::Relaxed)
91 }
92
93 fn is_subscribed(&self, topic: &str) -> bool {
94 self.subscriptions.all_topics().iter().any(|t| t == topic)
95 }
96
97 pub(super) async fn next(&mut self) -> Option<KrakenSpotWsMessage> {
99 if let Some(msg) = self.pending_messages.pop_front() {
100 return Some(msg);
101 }
102
103 loop {
104 tokio::select! {
105 Some(cmd) = self.cmd_rx.recv() => {
106 match cmd {
107 SpotHandlerCommand::SetClient(client) => {
108 log::debug!("WebSocketClient received by handler");
109 self.inner = Some(client);
110 }
111 SpotHandlerCommand::Disconnect => {
112 log::debug!("Disconnect command received");
113
114 if let Some(client) = self.inner.take() {
115 client.disconnect().await;
116 }
117 }
118 SpotHandlerCommand::Subscribe { payload }
119 | SpotHandlerCommand::Unsubscribe { payload } => {
120 if let Some(client) = &self.inner
121 && let Err(e) = client.send_text(payload.expose_secret().to_owned(), Some(KRAKEN_RATE_LIMIT_KEY_SUBSCRIPTION.as_slice())).await
122 {
123 log::error!("Failed to send text: {e}");
124 }
125 }
126 SpotHandlerCommand::Ping { payload } => {
127 if let Some(client) = &self.inner
128 && let Err(e) = client.send_text(payload.expose_secret().to_owned(), None).await
129 {
130 log::error!("Failed to send text: {e}");
131 }
132 }
133 SpotHandlerCommand::SendOrderRequest { req_id, payload } => {
134 if let Some(client) = &self.inner {
135 if let Err(e) = client.send_text(payload.expose_secret().to_owned(), Some(KRAKEN_RATE_LIMIT_KEY_ORDER.as_slice())).await {
136 log::error!(
137 "Kraken WS send_order_request failed req_id={req_id}: {e}"
138 );
139 } else {
140 log::debug!(
141 "Kraken WS send_order_request enqueued req_id={req_id}"
142 );
143 }
144 } else {
145 log::error!(
146 "Kraken WS send_order_request without active client req_id={req_id}"
147 );
148 }
149 }
150 }
151 }
152
153 msg = self.raw_rx.recv() => {
154 let msg = match msg {
155 Some(msg) => msg,
156 None => {
157 log::debug!("WebSocket stream closed");
158 return None;
159 }
160 };
161
162 if let Message::Ping(data) = &msg {
163 log::trace!("Received ping frame with {} bytes", data.len());
164
165 if let Some(client) = &self.inner
166 && let Err(e) = client.send_pong(data.to_vec()).await
167 {
168 log::warn!("Failed to send pong frame: {e}");
169 }
170 continue;
171 }
172
173 if self.signal.load(Ordering::Relaxed) {
174 log::debug!("Stop signal received");
175 return None;
176 }
177
178 let text = match msg {
179 Message::Text(text) => text.to_string(),
180 Message::Binary(data) => {
181 match String::from_utf8(data.to_vec()) {
182 Ok(text) => text,
183 Err(e) => {
184 log::warn!("Failed to decode binary message: {e}");
185 continue;
186 }
187 }
188 }
189 Message::Pong(_) => {
190 log::trace!("Received pong");
191 continue;
192 }
193 Message::Close(_) => {
194 log::debug!("WebSocket connection closed");
195 return None;
196 }
197 Message::Frame(_) => {
198 log::trace!("Received raw frame");
199 continue;
200 }
201 _ => continue,
202 };
203
204 if text == RECONNECTED {
205 log::debug!("Received WebSocket reconnected signal");
206 return Some(KrakenSpotWsMessage::Reconnected);
207 }
208
209 if let Some(msg) = self.parse_message(&text) {
210 return Some(msg);
211 }
212 }
213 }
214 }
215 }
216
217 fn parse_message(&self, text: &str) -> Option<KrakenSpotWsMessage> {
218 if text.len() < 50 && text.starts_with("{\"channel\":\"") {
220 if text.contains("heartbeat") {
221 log::trace!("Received heartbeat");
222 return None;
223 }
224
225 if text.contains("status") {
226 log::debug!("Received status message");
227 return None;
228 }
229 }
230
231 if text.contains("\"level3\"")
232 && let Some(msg) = parse_level3_text(text)
233 {
234 return self.handle_l3_message(msg);
235 }
236
237 let value: Value = match serde_json::from_str(text) {
238 Ok(v) => v,
239 Err(e) => {
240 log::warn!("Failed to parse message: {e}");
241 return None;
242 }
243 };
244
245 if value.get("method").is_some() {
246 match parse_order_response(text) {
247 Ok(Some(msg)) => return Some(msg),
248 Ok(None) => {}
249 Err(e) => log::warn!("Failed to parse order response: {e}"),
250 }
251 self.handle_control_message(value);
252 return None;
253 }
254
255 if value.get("channel").is_some() && value.get("data").is_some() {
256 match serde_json::from_str::<KrakenWsRawMessage>(text) {
257 Ok(msg) => return self.handle_data_message(msg),
258 Err(e) => {
259 log::debug!("Failed to parse data message: {e}");
260 return None;
261 }
262 }
263 }
264
265 log::debug!("Unhandled message structure: {text}");
266 None
267 }
268
269 fn handle_control_message(&self, value: Value) {
270 match serde_json::from_value::<KrakenWsResponse>(value) {
271 Ok(response) => match response {
272 KrakenWsResponse::Subscribe(sub) => {
273 if sub.success {
274 if let Some(result) = &sub.result {
275 log::debug!(
276 "Subscription confirmed: channel={:?}, req_id={:?}",
277 result.channel,
278 sub.req_id
279 );
280 } else {
281 log::debug!("Subscription confirmed: req_id={:?}", sub.req_id);
282 }
283 } else {
284 log::warn!(
285 "Subscription failed: error={:?}, req_id={:?}",
286 sub.error,
287 sub.req_id
288 );
289 }
290 }
291 KrakenWsResponse::Unsubscribe(unsub) => {
292 if unsub.success {
293 log::debug!("Unsubscription confirmed: req_id={:?}", unsub.req_id);
294 } else {
295 log::warn!(
296 "Unsubscription failed: error={:?}, req_id={:?}",
297 unsub.error,
298 unsub.req_id
299 );
300 }
301 }
302 KrakenWsResponse::Pong(pong) => {
303 log::trace!("Received pong: req_id={:?}", pong.req_id);
304 }
305 KrakenWsResponse::Other => {
306 log::debug!("Received unknown control response");
307 }
308 },
309 Err(_) => {
310 log::debug!("Received control message (failed to parse details)");
311 }
312 }
313 }
314
315 fn handle_data_message(&self, msg: KrakenWsRawMessage) -> Option<KrakenSpotWsMessage> {
316 match msg.channel {
317 KrakenWsChannel::Book => self.handle_book_message(msg),
318 KrakenWsChannel::Ticker => self.handle_ticker_message(msg),
319 KrakenWsChannel::Trade => self.handle_trade_message(msg),
320 KrakenWsChannel::Ohlc => self.handle_ohlc_message(msg),
321 KrakenWsChannel::Executions => self.handle_executions_message(msg),
322 KrakenWsChannel::Level3 => {
323 unreachable!("level3 messages routed via fast-path in parse_message",)
324 }
325 _ => {
326 log::warn!("Unhandled channel: {:?}", msg.channel);
327 None
328 }
329 }
330 }
331
332 fn handle_book_message(&self, msg: KrakenWsRawMessage) -> Option<KrakenSpotWsMessage> {
333 let is_snapshot = msg.event_type == KrakenWsMessageType::Snapshot;
334 let mut book_data = Vec::new();
335
336 for data in msg.data {
337 match serde_json::from_str::<KrakenWsBookData>(data.get()) {
338 Ok(bd) => {
339 if !self.is_subscribed(&format!("book:{}", bd.symbol)) {
340 continue;
341 }
342 book_data.push(bd);
343 }
344 Err(e) => log::error!("Failed to deserialize book data: {e}"),
345 }
346 }
347
348 if book_data.is_empty() {
349 None
350 } else {
351 Some(KrakenSpotWsMessage::Book {
352 data: book_data,
353 is_snapshot,
354 })
355 }
356 }
357
358 fn handle_ticker_message(&self, msg: KrakenWsRawMessage) -> Option<KrakenSpotWsMessage> {
359 let mut tickers = Vec::new();
360
361 for data in msg.data {
362 match serde_json::from_str::<KrakenWsTickerData>(data.get()) {
363 Ok(td) => {
364 let symbol = &td.symbol;
365 let quotes_key = format!("quotes:{symbol}");
366 let ticker_key = format!("ticker:{symbol}");
367 if !self.is_subscribed("es_key) && !self.is_subscribed(&ticker_key) {
368 continue;
369 }
370 tickers.push(td);
371 }
372 Err(e) => log::error!("Failed to deserialize ticker data: {e}"),
373 }
374 }
375
376 if tickers.is_empty() {
377 None
378 } else {
379 Some(KrakenSpotWsMessage::Ticker(tickers))
380 }
381 }
382
383 fn handle_trade_message(&self, msg: KrakenWsRawMessage) -> Option<KrakenSpotWsMessage> {
384 let mut trades = Vec::new();
385
386 for data in msg.data {
387 match serde_json::from_str::<KrakenWsTradeData>(data.get()) {
388 Ok(td) => trades.push(td),
389 Err(e) => log::error!("Failed to deserialize trade data: {e}"),
390 }
391 }
392
393 if trades.is_empty() {
394 None
395 } else {
396 Some(KrakenSpotWsMessage::Trade(trades))
397 }
398 }
399
400 fn handle_ohlc_message(&self, msg: KrakenWsRawMessage) -> Option<KrakenSpotWsMessage> {
401 let mut ohlc_data = Vec::new();
402
403 for data in msg.data {
404 match serde_json::from_str::<KrakenWsOhlcData>(data.get()) {
405 Ok(od) => ohlc_data.push(od),
406 Err(e) => log::error!("Failed to deserialize OHLC data: {e}"),
407 }
408 }
409
410 if ohlc_data.is_empty() {
411 None
412 } else {
413 Some(KrakenSpotWsMessage::Ohlc(ohlc_data))
414 }
415 }
416
417 fn handle_executions_message(&self, msg: KrakenWsRawMessage) -> Option<KrakenSpotWsMessage> {
418 let mut executions = Vec::new();
419
420 for data in msg.data {
421 match serde_json::from_str::<KrakenWsExecutionData>(data.get()) {
422 Ok(ed) => executions.push(ed),
423 Err(e) => log::error!("Failed to deserialize execution data: {e}"),
424 }
425 }
426
427 if executions.is_empty() {
428 None
429 } else {
430 Some(KrakenSpotWsMessage::Execution(executions))
431 }
432 }
433}
434
435#[derive(Deserialize)]
436struct Level3RawMessage<'a> {
437 channel: &'a str,
438 #[serde(rename = "type")]
439 msg_type: &'a str,
440 #[serde(borrow)]
441 data: Vec<&'a RawValue>,
442}
443
444fn parse_level3_text(text: &str) -> Option<KrakenSpotWsMessage> {
445 let msg: Level3RawMessage<'_> = serde_json::from_str(text).ok()?;
446 if msg.channel != "level3" {
447 return None;
448 }
449 let first = msg.data.first()?.get();
450
451 match msg.msg_type {
452 "snapshot" => match serde_json::from_str::<KrakenL3Snapshot>(first) {
453 Ok(snap) => Some(KrakenSpotWsMessage::L3Snapshot(snap)),
454 Err(e) => {
455 log::warn!("Failed to deserialize L3 snapshot: {e}");
456 None
457 }
458 },
459 "update" => match serde_json::from_str::<KrakenL3UpdateData>(first) {
460 Ok(update) => Some(KrakenSpotWsMessage::L3Update(update)),
461 Err(e) => {
462 log::warn!("Failed to deserialize L3 update: {e}");
463 None
464 }
465 },
466 _ => None,
467 }
468}
469
470impl SpotFeedHandler {
471 fn handle_l3_message(&self, msg: KrakenSpotWsMessage) -> Option<KrakenSpotWsMessage> {
472 let symbol = match &msg {
473 KrakenSpotWsMessage::L3Snapshot(s) => &s.symbol,
474 KrakenSpotWsMessage::L3Update(u) => &u.symbol,
475 _ => return None,
476 };
477
478 if !self.is_subscribed(&format!("level3:{symbol}")) {
479 return None;
480 }
481 Some(msg)
482 }
483}
484
485#[cfg(test)]
486mod tests {
487 use rstest::rstest;
488 use rust_decimal_macros::dec;
489
490 use super::*;
491
492 fn create_test_handler() -> SpotFeedHandler {
493 let signal = Arc::new(AtomicBool::new(false));
494 let (_cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
495 let (_raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel();
496 let subscriptions = SubscriptionState::new(':');
497
498 SpotFeedHandler::new(signal, cmd_rx, raw_rx, subscriptions)
499 }
500
501 #[rstest]
502 fn test_ticker_message_filtered_without_quotes_subscription() {
503 let handler = create_test_handler();
504
505 let json = r#"{
506 "channel": "ticker",
507 "type": "snapshot",
508 "data": [{
509 "symbol": "BTC/USD",
510 "bid": 105944.20,
511 "bid_qty": 2.5,
512 "ask": 105944.30,
513 "ask_qty": 3.2,
514 "last": 105899.40,
515 "volume": 163.28908096,
516 "vwap": 105904.39279,
517 "low": 104711.00,
518 "high": 106613.10,
519 "change": 250.00,
520 "change_pct": 0.24,
521 "timestamp": "2022-12-25T09:30:59.123456Z"
522 }]
523 }"#;
524
525 let result = handler.parse_message(json);
526 assert!(
527 result.is_none(),
528 "Ticker message should be filtered when no quotes subscription exists"
529 );
530 }
531
532 #[rstest]
533 fn test_ticker_message_passes_with_quotes_subscription() {
534 let handler = create_test_handler();
535 handler.subscriptions.mark_subscribe("quotes:BTC/USD");
536 handler.subscriptions.confirm_subscribe("quotes:BTC/USD");
537
538 let json = r#"{
539 "channel": "ticker",
540 "type": "snapshot",
541 "data": [{
542 "symbol": "BTC/USD",
543 "bid": 105944.20,
544 "bid_qty": 2.5,
545 "ask": 105944.30,
546 "ask_qty": 3.2,
547 "last": 105899.40,
548 "volume": 163.28908096,
549 "vwap": 105904.39279,
550 "low": 104711.00,
551 "high": 106613.10,
552 "change": 250.00,
553 "change_pct": 0.24,
554 "timestamp": "2022-12-25T09:30:59.123456Z"
555 }]
556 }"#;
557
558 let result = handler.parse_message(json);
559 assert!(
560 result.is_some(),
561 "Ticker message should pass with quotes subscription"
562 );
563
564 match result.unwrap() {
565 KrakenSpotWsMessage::Ticker(data) => {
566 assert!(!data.is_empty(), "Should have ticker data");
567 }
568 _ => panic!("Expected Ticker message"),
569 }
570 }
571
572 #[rstest]
573 fn test_ticker_message_passes_with_ticker_subscription() {
574 let handler = create_test_handler();
575 handler.subscriptions.mark_subscribe("ticker:BTC/USD");
576 handler.subscriptions.confirm_subscribe("ticker:BTC/USD");
577
578 let json = r#"{
579 "channel": "ticker",
580 "type": "snapshot",
581 "data": [{
582 "symbol": "BTC/USD",
583 "bid": 105944.20,
584 "bid_qty": 2.5,
585 "ask": 105944.30,
586 "ask_qty": 3.2,
587 "last": 105899.40,
588 "volume": 163.28908096,
589 "vwap": 105904.39279,
590 "low": 104711.00,
591 "high": 106613.10,
592 "change": 250.00,
593 "change_pct": 0.24,
594 "timestamp": "2022-12-25T09:30:59.123456Z"
595 }]
596 }"#;
597
598 let result = handler.parse_message(json);
599 assert!(
600 result.is_some(),
601 "Ticker message should pass with ticker: subscription"
602 );
603
604 match result.unwrap() {
605 KrakenSpotWsMessage::Ticker(data) => {
606 assert!(!data.is_empty(), "Should have ticker data");
607 }
608 _ => panic!("Expected Ticker message"),
609 }
610 }
611
612 #[rstest]
613 fn test_ticker_message_preserves_decimal_precision() {
614 let handler = create_test_handler();
615 handler.subscriptions.mark_subscribe("ticker:BTC/USD");
616 handler.subscriptions.confirm_subscribe("ticker:BTC/USD");
617 let json = include_str!("../../../test_data/ws_ticker_precision.json");
618
619 let message = handler.parse_message(json).unwrap();
620
621 let KrakenSpotWsMessage::Ticker(data) = message else {
622 panic!("Expected Ticker message, was {message:?}");
623 };
624 let ticker = &data[0];
625 assert_eq!(ticker.symbol, "BTC/USD");
626 assert_eq!(ticker.bid, dec!(123456789.123456789));
627 assert_eq!(ticker.bid_qty, dec!(0.1234567890123456789012345678));
628 assert_eq!(ticker.ask, dec!(123456789.223456789));
629 assert_eq!(ticker.ask_qty, dec!(0.2234567890123456789012345678));
630 assert_eq!(ticker.last, dec!(123456789.323456789));
631 assert_eq!(ticker.volume, dec!(123456789.423456789));
632 assert_eq!(ticker.vwap, dec!(123456789.523456789));
633 assert_eq!(ticker.low, dec!(123456789.623456789));
634 assert_eq!(ticker.high, dec!(123456789.723456789));
635 assert_eq!(ticker.change, dec!(123456789.823456789));
636 assert_eq!(ticker.change_pct, dec!(0.9234567890123456789012345678));
637 assert_eq!(
638 ticker.timestamp,
639 "2022-12-25T09:30:59.123456Z"
640 .parse::<jiff::Timestamp>()
641 .unwrap()
642 );
643 }
644
645 #[rstest]
646 fn test_book_message_filtered_without_book_subscription() {
647 let handler = create_test_handler();
648
649 let json = r#"{
650 "channel": "book",
651 "type": "snapshot",
652 "data": [{
653 "symbol": "BTC/USD",
654 "bids": [{"price": 105944.20, "qty": 2.5}],
655 "asks": [{"price": 105944.30, "qty": 3.2}],
656 "checksum": 12345,
657 "timestamp": "2023-10-06T17:35:55.440295Z"
658 }]
659 }"#;
660
661 let result = handler.parse_message(json);
662 assert!(
663 result.is_none(),
664 "Book message should be filtered when no book subscription exists"
665 );
666 }
667
668 #[rstest]
669 fn test_book_message_passes_with_book_subscription() {
670 let handler = create_test_handler();
671 handler.subscriptions.mark_subscribe("book:BTC/USD");
672 handler.subscriptions.confirm_subscribe("book:BTC/USD");
673
674 let json = r#"{
675 "channel": "book",
676 "type": "snapshot",
677 "data": [{
678 "symbol": "BTC/USD",
679 "bids": [{"price": 105944.20, "qty": 2.5}],
680 "asks": [{"price": 105944.30, "qty": 3.2}],
681 "checksum": 12345,
682 "timestamp": "2023-10-06T17:35:55.440295Z"
683 }]
684 }"#;
685
686 let result = handler.parse_message(json);
687 assert!(
688 result.is_some(),
689 "Book message should pass with book subscription"
690 );
691
692 match result.unwrap() {
693 KrakenSpotWsMessage::Book { data, is_snapshot } => {
694 assert!(!data.is_empty());
695 assert!(is_snapshot);
696 }
697 _ => panic!("Expected Book message"),
698 }
699 }
700
701 #[rstest]
702 fn test_send_order_request_variant_construction() {
703 let cmd = SpotHandlerCommand::SendOrderRequest {
704 req_id: 7,
705 payload: SecretString::from(r#"{"method":"add_order","req_id":7}"#.to_string()),
706 };
707
708 match cmd {
709 SpotHandlerCommand::SendOrderRequest { req_id, payload } => {
710 assert_eq!(req_id, 7);
711 assert!(payload.expose_secret().contains("add_order"));
712 }
713 _ => panic!("Expected SendOrderRequest, was a different variant"),
714 }
715 }
716
717 #[rstest]
718 fn test_command_debug_redacts_payload() {
719 let cmd = SpotHandlerCommand::SendOrderRequest {
720 req_id: 7,
721 payload: SecretString::from("auth-token".to_string()),
722 };
723
724 let debug = format!("{cmd:?}");
725
726 assert!(debug.contains(REDACTED));
727 assert!(!debug.contains("auth-token"));
728 }
729
730 #[rstest]
731 #[tokio::test]
732 async fn test_send_order_request_without_active_client_does_not_panic() {
733 let signal = Arc::new(AtomicBool::new(false));
734 let (cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel();
735 let (raw_tx, raw_rx) = tokio::sync::mpsc::unbounded_channel::<Message>();
736 let subscriptions = SubscriptionState::new(':');
737
738 let mut handler = SpotFeedHandler::new(signal.clone(), cmd_rx, raw_rx, subscriptions);
739
740 cmd_tx
741 .send(SpotHandlerCommand::SendOrderRequest {
742 req_id: 42,
743 payload: SecretString::from(r#"{"method":"add_order","req_id":42}"#.to_string()),
744 })
745 .unwrap();
746
747 drop(cmd_tx);
748 drop(raw_tx);
749
750 let result = handler.next().await;
751 assert!(
752 result.is_none(),
753 "Handler should return None when streams close"
754 );
755 }
756
757 #[rstest]
758 fn test_quotes_and_book_subscriptions_independent() {
759 let handler = create_test_handler();
760 handler.subscriptions.mark_subscribe("quotes:BTC/USD");
761 handler.subscriptions.confirm_subscribe("quotes:BTC/USD");
762
763 let book_json = r#"{
764 "channel": "book",
765 "type": "snapshot",
766 "data": [{
767 "symbol": "BTC/USD",
768 "bids": [{"price": 105944.20, "qty": 2.5}],
769 "asks": [{"price": 105944.30, "qty": 3.2}],
770 "checksum": 12345,
771 "timestamp": "2023-10-06T17:35:55.440295Z"
772 }]
773 }"#;
774
775 let book_result = handler.parse_message(book_json);
776 assert!(
777 book_result.is_none(),
778 "Book message should be filtered without book: subscription"
779 );
780
781 let ticker_json = r#"{
782 "channel": "ticker",
783 "type": "snapshot",
784 "data": [{
785 "symbol": "BTC/USD",
786 "bid": 105944.20,
787 "bid_qty": 2.5,
788 "ask": 105944.30,
789 "ask_qty": 3.2,
790 "last": 105899.40,
791 "volume": 163.28908096,
792 "vwap": 105904.39279,
793 "low": 104711.00,
794 "high": 106613.10,
795 "change": 250.00,
796 "change_pct": 0.24,
797 "timestamp": "2022-12-25T09:30:59.123456Z"
798 }]
799 }"#;
800
801 let ticker_result = handler.parse_message(ticker_json);
802 assert!(
803 ticker_result.is_some(),
804 "Ticker should pass with quotes subscription"
805 );
806 }
807
808 #[rstest]
809 fn test_parse_message_routes_add_order_response_to_order_response_variant() {
810 use super::super::enums::KrakenWsMethod;
811
812 let handler = create_test_handler();
813 let json = r#"{"method":"add_order","req_id":42,"success":true,"time_in":"2026-05-05T10:00:00.123Z","time_out":"2026-05-05T10:00:00.125Z","result":{"order_id":"OABCDE-12345-FGHIJ","cl_ord_id":"O-20260505-000001","order_userref":0}}"#;
814
815 let result = handler.parse_message(json);
816 match result {
817 Some(KrakenSpotWsMessage::OrderResponse(resp)) => {
818 assert_eq!(resp.method, KrakenWsMethod::AddOrder);
819 assert_eq!(resp.req_id, Some(42));
820 assert!(resp.success);
821 }
822 other => panic!("expected OrderResponse, was {other:?}"),
823 }
824 }
825
826 #[rstest]
827 fn test_parse_level3_snapshot_with_subscription_passes() {
828 let handler = create_test_handler();
829 handler.subscriptions.mark_subscribe("level3:BTC/USD");
830 handler.subscriptions.confirm_subscribe("level3:BTC/USD");
831
832 let json = r#"{
833 "channel": "level3",
834 "type": "snapshot",
835 "data": [{
836 "symbol": "BTC/USD",
837 "bids": [],
838 "asks": [],
839 "checksum": 0,
840 "timestamp": "2024-01-01T00:00:00Z"
841 }]
842 }"#;
843
844 let result = handler.parse_message(json);
845 assert!(matches!(result, Some(KrakenSpotWsMessage::L3Snapshot(_))));
846 }
847
848 #[rstest]
849 fn test_parse_level3_update_without_subscription_filtered() {
850 let handler = create_test_handler();
851 let json = r#"{
852 "channel": "level3",
853 "type": "update",
854 "data": [{
855 "symbol": "BTC/USD",
856 "bids": [],
857 "asks": [],
858 "checksum": 0,
859 "timestamp": "2024-01-01T00:00:00Z"
860 }]
861 }"#;
862
863 let result = handler.parse_message(json);
864 assert!(result.is_none());
865 }
866
867 #[rstest]
868 fn test_parse_level3_snapshot_compact_json() {
869 let handler = create_test_handler();
870 handler.subscriptions.mark_subscribe("level3:BTC/USD");
871 handler.subscriptions.confirm_subscribe("level3:BTC/USD");
872 let json = r#"{"channel":"level3","type":"snapshot","data":[{"symbol":"BTC/USD","bids":[],"asks":[],"checksum":0,"timestamp":"2024-01-01T00:00:00Z"}]}"#;
873 assert!(matches!(
874 handler.parse_message(json),
875 Some(KrakenSpotWsMessage::L3Snapshot(_))
876 ));
877 }
878
879 #[rstest]
880 fn test_parse_level3_snapshot_preserves_raw_decimal() {
881 let handler = create_test_handler();
882 handler.subscriptions.mark_subscribe("level3:BTC/USD");
883 handler.subscriptions.confirm_subscribe("level3:BTC/USD");
884
885 let json = r#"{
886 "channel": "level3",
887 "type": "snapshot",
888 "data": [{
889 "symbol": "BTC/USD",
890 "bids": [{
891 "order_id": "order-bid-1",
892 "limit_price": 42000.50000,
893 "order_qty": 0.01000000,
894 "timestamp": "2024-01-01T00:00:00Z"
895 }],
896 "asks": [],
897 "checksum": 0,
898 "timestamp": "2024-01-01T00:00:00Z"
899 }]
900 }"#;
901
902 let Some(KrakenSpotWsMessage::L3Snapshot(snap)) = handler.parse_message(json) else {
903 panic!("expected L3 snapshot");
904 };
905 assert_eq!(snap.bids[0].limit_price.raw, "42000.50000");
906 assert_eq!(snap.bids[0].order_qty.raw, "0.01000000");
907 }
908}