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