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