Skip to main content

nautilus_kraken/websocket/spot_v2/
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 Spot v2.
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, 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/// Commands sent from the outer client to the inner message handler.
48#[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
58/// WebSocket message handler for Kraken Spot v2.
59pub(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    /// Creates a new [`SpotFeedHandler`] instance.
70    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    /// Processes messages and commands, returning when stopped or stream ends.
95    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        // Fast pre-filter for high-frequency control messages (no JSON parsing)
216        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(&quotes_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}