Skip to main content

nautilus_kraken/data/
spot.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//! Kraken Spot data client implementation.
17
18use std::{
19    future::Future,
20    sync::{
21        Arc,
22        atomic::{AtomicBool, AtomicU64, Ordering},
23    },
24    time::Duration,
25};
26
27use ahash::AHashMap;
28use anyhow::Context;
29use async_trait::async_trait;
30use futures_util::StreamExt;
31use nautilus_common::{
32    clients::DataClient,
33    live::get_data_event_sender,
34    messages::{
35        DataEvent,
36        data::{
37            BarsResponse, BookResponse, DataResponse, InstrumentResponse, InstrumentsResponse,
38            RequestBars, RequestBookSnapshot, RequestInstrument, RequestInstruments, RequestTrades,
39            SubscribeBars, SubscribeBookDeltas, SubscribeIndexPrices, SubscribeInstrument,
40            SubscribeInstrumentStatus, SubscribeInstruments, SubscribeMarkPrices, SubscribeQuotes,
41            SubscribeTrades, TradesResponse, UnsubscribeBars, UnsubscribeBookDeltas,
42            UnsubscribeIndexPrices, UnsubscribeInstrumentStatus, UnsubscribeMarkPrices,
43            UnsubscribeQuotes, UnsubscribeTrades,
44        },
45    },
46};
47use nautilus_core::{
48    AtomicMap, UnixNanos,
49    datetime::datetime_to_unix_nanos,
50    time::{AtomicTime, get_atomic_clock_realtime},
51};
52use nautilus_live::{SocketControlFactory, task::TaskGroup};
53use nautilus_model::{
54    data::{Bar, Data, OrderBookDeltas},
55    enums::{AggregationSource, BookType},
56    identifiers::{ClientId, InstrumentId, Venue},
57    instruments::{Instrument, InstrumentAny},
58};
59use parking_lot::Mutex;
60use tokio_util::sync::CancellationToken;
61use ustr::Ustr;
62
63type OhlcBufferKey = (Ustr, u32);
64type OhlcBuffer = Arc<Mutex<AHashMap<OhlcBufferKey, (Bar, UnixNanos)>>>;
65
66use crate::{
67    common::{consts::KRAKEN_VENUE, lookup_instrument_in_snapshot},
68    config::KrakenDataClientConfig,
69    http::{KrakenSpotHttpClient, spot::client::KRAKEN_SPOT_DEFAULT_RATE_LIMIT_PER_SECOND},
70    websocket::spot_v2::{
71        client::KrakenSpotWebSocketClient,
72        level_2::{L2BookState, L2Depths},
73        level_3::{
74            BookOrderIdHasher, KrakenL3WsMessage,
75            resync::retry_l3_resync,
76            runtime::{L3Sink, L3State, process_l3_message},
77        },
78        messages::KrakenSpotWsMessage,
79        parse::{parse_quote_tick, parse_trade_tick, parse_ws_bar},
80    },
81};
82
83/// `L3Sink` implementation that forwards deltas to the data engine.
84struct DataEventSink<'a> {
85    sender: &'a tokio::sync::mpsc::UnboundedSender<DataEvent>,
86}
87
88impl L3Sink for DataEventSink<'_> {
89    fn emit_deltas(&mut self, deltas: OrderBookDeltas) {
90        if let Err(e) = self
91            .sender
92            .send(DataEvent::Data(Data::Deltas(Box::new(deltas))))
93        {
94            log::error!("Failed to send L3 deltas: {e}");
95        }
96    }
97}
98
99struct SpotMessageContext<'a> {
100    sender: &'a tokio::sync::mpsc::UnboundedSender<DataEvent>,
101    instruments: &'a Arc<AtomicMap<InstrumentId, InstrumentAny>>,
102    book_sequence: &'a Arc<AtomicU64>,
103    l2_depths: &'a L2Depths,
104    ohlc_buffer: &'a OhlcBuffer,
105    clock: &'static AtomicTime,
106}
107
108/// Kraken Spot data client.
109///
110/// Provides real-time market data from Kraken Spot markets through WebSocket v2.
111#[allow(dead_code)]
112#[derive(Debug)]
113pub struct KrakenSpotDataClient {
114    clock: &'static AtomicTime,
115    client_id: ClientId,
116    config: KrakenDataClientConfig,
117    http: KrakenSpotHttpClient,
118    ws: KrakenSpotWebSocketClient,
119    ws_l3: Option<KrakenSpotWebSocketClient>,
120    socket_factory: SocketControlFactory,
121    l3_handler_alive: Arc<AtomicBool>,
122    is_connected: AtomicBool,
123    cancellation_token: CancellationToken,
124    session_tasks: TaskGroup,
125    command_tasks: TaskGroup,
126    instruments: Arc<AtomicMap<InstrumentId, InstrumentAny>>,
127    data_sender: tokio::sync::mpsc::UnboundedSender<DataEvent>,
128}
129
130impl KrakenSpotDataClient {
131    /// Creates a new [`KrakenSpotDataClient`] instance.
132    pub fn new(client_id: ClientId, config: KrakenDataClientConfig) -> anyhow::Result<Self> {
133        let session_tasks = TaskGroup::new();
134        let cancellation_token = session_tasks.cancellation_token();
135        let command_tasks = TaskGroup::new();
136        let socket_factory = SocketControlFactory::new(client_id, Some(*KRAKEN_VENUE));
137
138        let http = KrakenSpotHttpClient::new(
139            config.environment,
140            config.base_url.clone(),
141            config.timeout_secs,
142            None,
143            None,
144            None,
145            config.proxy_url.clone(),
146            config
147                .max_requests_per_second
148                .unwrap_or(KRAKEN_SPOT_DEFAULT_RATE_LIMIT_PER_SECOND),
149        )?;
150
151        let ws = KrakenSpotWebSocketClient::new(
152            config.clone(),
153            cancellation_token.clone(),
154            config.proxy_url.clone(),
155        )
156        .with_socket_control(socket_factory.control("kraken-spot-data-streams"));
157
158        Ok(Self {
159            clock: get_atomic_clock_realtime(),
160            client_id,
161            config,
162            http,
163            ws,
164            ws_l3: None,
165            socket_factory,
166            l3_handler_alive: Arc::new(AtomicBool::new(false)),
167            is_connected: AtomicBool::new(false),
168            cancellation_token,
169            session_tasks,
170            command_tasks,
171            instruments: Arc::new(AtomicMap::new()),
172            data_sender: get_data_event_sender(),
173        })
174    }
175
176    /// Returns the cached instruments.
177    #[must_use]
178    pub fn instruments(&self) -> Vec<InstrumentAny> {
179        self.instruments.load().values().cloned().collect()
180    }
181
182    /// Returns a cached instrument by ID.
183    #[must_use]
184    pub fn get_instrument(&self, instrument_id: &InstrumentId) -> Option<InstrumentAny> {
185        self.instruments.load().get(instrument_id).cloned()
186    }
187
188    async fn load_instruments(&self) -> anyhow::Result<Vec<InstrumentAny>> {
189        let instruments = self
190            .http
191            .request_instruments(None)
192            .await
193            .context("Failed to load spot instruments")?;
194
195        self.instruments.rcu(|m| {
196            for instrument in &instruments {
197                m.insert(instrument.id(), instrument.clone());
198            }
199        });
200
201        self.http.cache_instruments(&instruments);
202
203        log::debug!(
204            "Loaded instruments: client_id={}, count={}",
205            self.client_id,
206            instruments.len()
207        );
208
209        Ok(instruments)
210    }
211
212    fn spawn_ws<F>(&self, fut: F, context: &'static str)
213    where
214        F: Future<Output = anyhow::Result<()>> + Send + 'static,
215    {
216        let future = async move {
217            if let Err(e) = fut.await {
218                log::error!("{context}: {e:?}");
219            }
220        };
221
222        if let Err(e) = self.command_tasks.spawn(future) {
223            log::warn!("Skipping Kraken Spot {context} after shutdown began: {e}");
224        }
225    }
226
227    fn spawn_command<F>(&self, future: F)
228    where
229        F: Future<Output = ()> + Send + 'static,
230    {
231        if let Err(e) = self.command_tasks.spawn(future) {
232            log::warn!("Skipping Kraken Spot data command after shutdown began: {e}");
233        }
234    }
235
236    async fn finish_tasks(&self) -> anyhow::Result<()> {
237        let (session_result, command_result) = tokio::join!(
238            self.session_tasks
239                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
240            self.command_tasks
241                .finish_shutdown(Duration::from_secs(1), Duration::from_secs(2)),
242        );
243        session_result.context("failed to finish Kraken Spot data session tasks")?;
244        command_result.context("failed to finish Kraken Spot data command tasks")?;
245        Ok(())
246    }
247
248    async fn prepare_task_groups(&mut self) -> anyhow::Result<()> {
249        if !self.session_tasks.is_open() || !self.command_tasks.is_open() {
250            self.session_tasks.begin_shutdown();
251            self.command_tasks.begin_shutdown();
252            self.ws
253                .close()
254                .await
255                .context("failed to close prior Kraken Spot WebSocket")?;
256
257            if let Some(ws_l3) = self.ws_l3.as_mut() {
258                ws_l3
259                    .close()
260                    .await
261                    .context("failed to close prior Kraken Spot L3 WebSocket")?;
262                self.ws_l3 = None;
263            }
264            self.finish_tasks().await?;
265            self.session_tasks
266                .start_generation()
267                .context("failed to start Kraken Spot data session task generation")?;
268            self.command_tasks
269                .start_generation()
270                .context("failed to start Kraken Spot data command task generation")?;
271            self.cancellation_token = self.session_tasks.cancellation_token();
272            self.ws = KrakenSpotWebSocketClient::new(
273                self.config.clone(),
274                self.cancellation_token.clone(),
275                self.config.proxy_url.clone(),
276            )
277            .with_socket_control(self.socket_factory.control("kraken-spot-data-streams"));
278        }
279        Ok(())
280    }
281
282    async fn teardown_partial_connect(&mut self) -> anyhow::Result<()> {
283        self.session_tasks.begin_shutdown();
284        self.command_tasks.begin_shutdown();
285        let ws_result = self.ws.close().await;
286        let ws_l3_result = if let Some(ws_l3) = self.ws_l3.as_mut() {
287            ws_l3.close().await
288        } else {
289            Ok(())
290        };
291
292        if ws_l3_result.is_ok() {
293            self.ws_l3 = None;
294        }
295        let tasks_result = self.finish_tasks().await;
296        self.is_connected.store(false, Ordering::Release);
297        tasks_result?;
298        ws_result?;
299        Ok(ws_l3_result?)
300    }
301
302    fn subscribe_l3_book(&mut self, cmd: &SubscribeBookDeltas) -> anyhow::Result<()> {
303        let instrument_id = cmd.instrument_id;
304        let symbol_ustr = instrument_id.symbol.inner();
305        let depth = cmd.depth.map_or(1000, |d| d.get() as u32);
306
307        if !matches!(depth, 10 | 100 | 1000) {
308            anyhow::bail!("Invalid L3 depth {depth} for Kraken Spot, valid values: 10, 100, 1000");
309        }
310
311        if !self.config.has_api_credentials() {
312            anyhow::bail!(
313                "L3 order book requires API credentials; configure api_key and api_secret"
314            );
315        }
316
317        let handler_dead = !self.l3_handler_alive.load(Ordering::Relaxed);
318
319        if self.ws_l3.is_none() {
320            let ws_l3 = KrakenSpotWebSocketClient::l3(
321                self.config.clone(),
322                self.cancellation_token.clone(),
323                self.config.proxy_url.clone(),
324            )
325            .with_socket_control(self.socket_factory.control("kraken-spot-l3-data-streams"));
326
327            self.spawn_l3_handler_task(ws_l3.clone(), false);
328            self.ws_l3 = Some(ws_l3);
329        } else if handler_dead && let Some(ws_l3) = self.ws_l3.as_ref() {
330            let ws_l3 = ws_l3.clone();
331            self.spawn_l3_handler_task(ws_l3, true);
332        }
333
334        let ws_l3 = self
335            .ws_l3
336            .as_ref()
337            .expect("ws_l3 initialised above")
338            .clone();
339
340        self.spawn_ws(
341            async move {
342                ws_l3
343                    .wait_until_active(10.0)
344                    .await
345                    .map_err(|e| anyhow::anyhow!("L3 WebSocket failed to become active: {e}"))?;
346                ws_l3
347                    .wait_until_authenticated(10.0)
348                    .await
349                    .map_err(|e| anyhow::anyhow!("L3 WebSocket failed to authenticate: {e}"))?;
350                ws_l3
351                    .subscribe_book_l3(symbol_ustr, depth)
352                    .await
353                    .map_err(|e| anyhow::anyhow!("{e}"))
354            },
355            "subscribe l3 book",
356        );
357
358        Ok(())
359    }
360
361    fn spawn_l3_handler_task(&self, handler_client: KrakenSpotWebSocketClient, restart: bool) {
362        let data_sender = self.data_sender.clone();
363        let instruments = self.instruments.clone();
364        let cancellation_token = self.cancellation_token.clone();
365        let clock = self.clock;
366        let alive = self.l3_handler_alive.clone();
367        let session_spawner = match self.session_tasks.spawner() {
368            Ok(spawner) => spawner,
369            Err(e) => {
370                log::warn!("Skipping Kraken L3 handler after shutdown began: {e}");
371                return;
372            }
373        };
374
375        alive.store(true, Ordering::Relaxed);
376
377        let future = async move {
378            struct AliveGuard(Arc<AtomicBool>);
379            impl Drop for AliveGuard {
380                fn drop(&mut self) {
381                    self.0.store(false, Ordering::Relaxed);
382                }
383            }
384            let _alive_guard = AliveGuard(alive);
385
386            let mut handler_client = handler_client;
387
388            if restart && let Err(e) = handler_client.close().await {
389                log::error!("Failed to close prior L3 WebSocket generation: {e}");
390                return;
391            }
392
393            if let Err(e) = handler_client.connect().await {
394                log::error!("L3 WebSocket connect failed: {e}");
395                return;
396            }
397
398            if let Err(e) = handler_client.wait_until_active(10.0).await {
399                log::error!("L3 WebSocket failed to become active: {e}");
400                return;
401            }
402
403            if let Err(e) = handler_client.authenticate().await {
404                log::error!("L3 WebSocket authentication failed: {e}");
405                return;
406            }
407
408            let stream = match handler_client.stream() {
409                Ok(s) => s,
410                Err(e) => {
411                    log::error!("L3 stream() failed: {e}");
412                    return;
413                }
414            };
415            tokio::pin!(stream);
416
417            let mut states: AHashMap<String, L3State> = AHashMap::new();
418            let hasher = BookOrderIdHasher::new();
419            let l3_depths = handler_client.l3_depths_handle();
420            let validate_checksum = handler_client.validate_l3_checksum();
421            let resync_client = handler_client.clone();
422
423            loop {
424                tokio::select! {
425                    () = cancellation_token.cancelled() => break,
426                    msg = stream.next() => {
427                        let Some(msg) = msg else { break };
428                        let ts_init = clock.get_time_ns();
429
430                        let runtime_msg = match msg {
431                            KrakenSpotWsMessage::L3Snapshot(snap) => {
432                                KrakenL3WsMessage::Snapshot(snap)
433                            }
434                            KrakenSpotWsMessage::L3Update(update) => {
435                                KrakenL3WsMessage::Update(update)
436                            }
437                            KrakenSpotWsMessage::Reconnected => {
438                                log::info!("L3 WebSocket reconnected");
439
440                                for state in states.values_mut() {
441                                    state.open_orders.clear();
442                                    state.awaiting_snapshot = true;
443                                }
444                                continue;
445                            }
446                            _ => continue,
447                        };
448
449                        let mut sink = DataEventSink { sender: &data_sender };
450                        let resync = process_l3_message(
451                            runtime_msg,
452                            &mut sink,
453                            &instruments,
454                            &l3_depths,
455                            &mut states,
456                            &hasher,
457                            validate_checksum,
458                            ts_init,
459                        );
460
461                        if let Some(request) = resync {
462                            log::warn!(
463                                "Resyncing Kraken L3 book: symbol={}, depth={}, reason={}",
464                                request.symbol,
465                                request.depth,
466                                request.reason,
467                            );
468                            let symbol_ustr = Ustr::from(&request.symbol);
469                            let client_for_resync = resync_client.clone();
470
471                            if let Err(e) = session_spawner.spawn(async move {
472                                retry_l3_resync(
473                                    &client_for_resync,
474                                    symbol_ustr,
475                                    request.depth,
476                                )
477                                .await;
478                            }) {
479                                log::warn!("Skipping Kraken L3 resync after shutdown began: {e}");
480                            }
481                        }
482                    }
483                }
484            }
485        };
486
487        if let Err(e) = self.session_tasks.spawn(future) {
488            self.l3_handler_alive.store(false, Ordering::Relaxed);
489            log::warn!("Skipping Kraken L3 handler after shutdown began: {e}");
490        }
491    }
492
493    fn spawn_message_handler(&mut self) -> anyhow::Result<()> {
494        let stream = self.ws.stream().map_err(|e| anyhow::anyhow!("{e}"))?;
495        let data_sender = self.data_sender.clone();
496        let instruments = self.instruments.clone();
497        let book_sequence = Arc::new(AtomicU64::new(0));
498        let ohlc_buffer: OhlcBuffer = Arc::new(Mutex::new(AHashMap::new()));
499        let l2_depths = self.ws.l2_depths_handle();
500        let cancellation_token = self.cancellation_token.clone();
501        let clock = self.clock;
502
503        let future = async move {
504            tokio::pin!(stream);
505            let mut l2_books = L2BookState::default();
506
507            loop {
508                tokio::select! {
509                    () = cancellation_token.cancelled() => {
510                        log::debug!("Spot message handler cancelled");
511                        Self::flush_ohlc_buffer(&ohlc_buffer, &data_sender);
512                        break;
513                    }
514                    msg = stream.next() => {
515                        match msg {
516                            Some(ws_msg) => {
517                                let context = SpotMessageContext {
518                                    sender: &data_sender,
519                                    instruments: &instruments,
520                                    book_sequence: &book_sequence,
521                                    l2_depths: &l2_depths,
522                                    ohlc_buffer: &ohlc_buffer,
523                                    clock,
524                                };
525                                Self::handle_ws_message(ws_msg, &context, &mut l2_books);
526                            }
527                            None => {
528                                log::debug!("Spot WebSocket stream ended");
529                                Self::flush_ohlc_buffer(&ohlc_buffer, &data_sender);
530                                break;
531                            }
532                        }
533                    }
534                }
535            }
536        };
537
538        self.session_tasks
539            .spawn(future)
540            .context("failed to register Kraken Spot message handler")
541    }
542
543    fn flush_ohlc_buffer(
544        ohlc_buffer: &OhlcBuffer,
545        sender: &tokio::sync::mpsc::UnboundedSender<DataEvent>,
546    ) {
547        let mut buffer = ohlc_buffer.lock();
548        let bars: Vec<Bar> = buffer.drain().map(|(_, (bar, _))| bar).collect();
549        for bar in bars {
550            if let Err(e) = sender.send(DataEvent::Data(Data::Bar(bar))) {
551                log::error!("Failed to send buffered bar: {e}");
552            }
553        }
554    }
555
556    fn handle_ws_message(
557        msg: KrakenSpotWsMessage,
558        context: &SpotMessageContext,
559        l2_books: &mut L2BookState,
560    ) {
561        let ts_init = context.clock.get_time_ns();
562
563        match msg {
564            KrakenSpotWsMessage::Ticker(tickers) => {
565                let instruments = context.instruments.load();
566
567                for ticker in &tickers {
568                    let Some(instrument) =
569                        lookup_instrument_in_snapshot(&instruments, ticker.symbol.as_str())
570                    else {
571                        log::warn!("No instrument for symbol: {}", ticker.symbol);
572                        continue;
573                    };
574
575                    match parse_quote_tick(ticker, instrument, ts_init) {
576                        Ok(quote) => {
577                            if let Err(e) = context.sender.send(DataEvent::Data(Data::Quote(quote)))
578                            {
579                                log::error!("Failed to send quote: {e}");
580                            }
581                        }
582                        Err(e) => log::error!("Failed to parse quote tick: {e}"),
583                    }
584                }
585            }
586            KrakenSpotWsMessage::Trade(trades) => {
587                let instruments = context.instruments.load();
588
589                for trade in &trades {
590                    let Some(instrument) =
591                        lookup_instrument_in_snapshot(&instruments, trade.symbol.as_str())
592                    else {
593                        log::warn!("No instrument for symbol: {}", trade.symbol);
594                        continue;
595                    };
596
597                    match parse_trade_tick(trade, instrument, ts_init) {
598                        Ok(tick) => {
599                            if let Err(e) = context.sender.send(DataEvent::Data(Data::Trade(tick)))
600                            {
601                                log::error!("Failed to send trade: {e}");
602                            }
603                        }
604                        Err(e) => log::error!("Failed to parse trade tick: {e}"),
605                    }
606                }
607            }
608            KrakenSpotWsMessage::Book { data, is_snapshot } => {
609                let instruments = context.instruments.load();
610
611                for book in &data {
612                    let Some(instrument) =
613                        lookup_instrument_in_snapshot(&instruments, book.symbol.as_str())
614                    else {
615                        log::warn!("No instrument for symbol: {}", book.symbol);
616                        continue;
617                    };
618                    let sequence = context.book_sequence.load(Ordering::Relaxed);
619                    let depth = context.l2_depths.get(book.symbol.as_str());
620                    match l2_books.process_book(
621                        book,
622                        instrument,
623                        sequence,
624                        is_snapshot,
625                        depth,
626                        ts_init,
627                    ) {
628                        Ok(Some((deltas, next_sequence))) => {
629                            context
630                                .book_sequence
631                                .store(next_sequence, Ordering::Relaxed);
632
633                            if let Err(e) = context
634                                .sender
635                                .send(DataEvent::Data(Data::Deltas(Box::new(deltas))))
636                            {
637                                log::error!("Failed to send deltas: {e}");
638                            }
639                        }
640                        Ok(None) => {}
641                        Err(e) => log::error!("Failed to parse book deltas: {e}"),
642                    }
643                }
644            }
645            KrakenSpotWsMessage::Ohlc(ohlc_data) => {
646                let mut buffer = context.ohlc_buffer.lock();
647
648                let instruments = context.instruments.load();
649
650                for ohlc in &ohlc_data {
651                    let Some(instrument) =
652                        lookup_instrument_in_snapshot(&instruments, ohlc.symbol.as_str())
653                    else {
654                        log::warn!("No instrument for symbol: {}", ohlc.symbol);
655                        continue;
656                    };
657
658                    match parse_ws_bar(ohlc, instrument, ts_init) {
659                        Ok(new_bar) => {
660                            let key: (Ustr, u32) = (ohlc.symbol, ohlc.interval);
661                            let new_interval_begin = UnixNanos::from(
662                                u64::try_from(ohlc.interval_begin.as_nanosecond()).unwrap_or(0),
663                            );
664
665                            if let Some((buffered_bar, buffered_begin)) = buffer.get(&key)
666                                && new_interval_begin != *buffered_begin
667                                && let Err(e) = context
668                                    .sender
669                                    .send(DataEvent::Data(Data::Bar(*buffered_bar)))
670                            {
671                                log::error!("Failed to send bar: {e}");
672                            }
673
674                            buffer.insert(key, (new_bar, new_interval_begin));
675                        }
676                        Err(e) => log::error!("Failed to parse bar: {e}"),
677                    }
678                }
679            }
680            KrakenSpotWsMessage::Execution(_) => {}
681            KrakenSpotWsMessage::OrderResponse(_) => {}
682            KrakenSpotWsMessage::L3Snapshot(_) => {}
683            KrakenSpotWsMessage::L3Update(_) => {}
684            KrakenSpotWsMessage::Reconnected => {
685                log::info!("Spot WebSocket reconnected");
686            }
687        }
688    }
689}
690
691#[async_trait(?Send)]
692impl DataClient for KrakenSpotDataClient {
693    fn client_id(&self) -> ClientId {
694        self.client_id
695    }
696
697    fn venue(&self) -> Option<Venue> {
698        Some(*KRAKEN_VENUE)
699    }
700
701    fn start(&mut self) -> anyhow::Result<()> {
702        log::info!(
703            "Starting Spot data client: client_id={}, environment={:?}",
704            self.client_id,
705            self.config.environment
706        );
707        Ok(())
708    }
709
710    fn stop(&mut self) -> anyhow::Result<()> {
711        log::info!("Stopping Spot data client: {}", self.client_id);
712        self.session_tasks.begin_shutdown();
713        self.command_tasks.begin_shutdown();
714        self.ws.begin_shutdown();
715        self.is_connected.store(false, Ordering::Relaxed);
716        Ok(())
717    }
718
719    fn reset(&mut self) -> anyhow::Result<()> {
720        log::info!("Resetting Spot data client: {}", self.client_id);
721        self.session_tasks.begin_shutdown();
722        self.command_tasks.begin_shutdown();
723        self.ws.begin_shutdown();
724        self.is_connected.store(false, Ordering::Relaxed);
725
726        self.instruments.store(ahash::AHashMap::new());
727        Ok(())
728    }
729
730    fn dispose(&mut self) -> anyhow::Result<()> {
731        log::debug!("Disposing Spot data client: {}", self.client_id);
732        self.stop()
733    }
734
735    fn is_connected(&self) -> bool {
736        self.is_connected.load(Ordering::SeqCst)
737    }
738
739    fn is_disconnected(&self) -> bool {
740        !self.is_connected()
741    }
742
743    async fn connect(&mut self) -> anyhow::Result<()> {
744        if self.is_connected() && self.session_tasks.is_open() && self.command_tasks.is_open() {
745            return Ok(());
746        }
747
748        self.prepare_task_groups().await?;
749
750        let instruments = self.load_instruments().await?;
751
752        let session_result = async {
753            self.ws
754                .connect()
755                .await
756                .context("Failed to connect spot WebSocket")?;
757            self.ws
758                .wait_until_active(10.0)
759                .await
760                .context("Spot WebSocket failed to become active")?;
761
762            self.spawn_message_handler()?;
763
764            Ok::<(), anyhow::Error>(())
765        }
766        .await;
767
768        if let Err(e) = session_result {
769            if let Err(teardown_error) = self.teardown_partial_connect().await {
770                return Err(e.context(format!(
771                    "Kraken Spot data startup teardown failed: {teardown_error}"
772                )));
773            }
774            return Err(e);
775        }
776
777        for instrument in instruments {
778            if let Err(e) = self.data_sender.send(DataEvent::Instrument(instrument)) {
779                log::error!("Failed to send instrument: {e}");
780            }
781        }
782
783        self.is_connected.store(true, Ordering::Release);
784        log::info!("Connected: client_id={}, product_type=Spot", self.client_id);
785        Ok(())
786    }
787
788    async fn disconnect(&mut self) -> anyhow::Result<()> {
789        self.teardown_partial_connect().await?;
790        self.is_connected.store(false, Ordering::Relaxed);
791
792        log::info!("Disconnected: client_id={}", self.client_id);
793        Ok(())
794    }
795
796    fn subscribe_instruments(&mut self, _cmd: SubscribeInstruments) -> anyhow::Result<()> {
797        log::debug!("subscribe_instruments: Kraken instruments are fetched via HTTP on connect");
798        Ok(())
799    }
800
801    fn subscribe_instrument(&mut self, _cmd: SubscribeInstrument) -> anyhow::Result<()> {
802        log::debug!("subscribe_instrument: Kraken instruments are fetched via HTTP on connect");
803        Ok(())
804    }
805
806    fn subscribe_book_deltas(&mut self, cmd: SubscribeBookDeltas) -> anyhow::Result<()> {
807        let instrument_id = cmd.instrument_id;
808        let depth = cmd.depth;
809
810        match cmd.book_type {
811            BookType::L2_MBP => {}
812            BookType::L3_MBO => return self.subscribe_l3_book(&cmd),
813            other => {
814                log::warn!("Unsupported BookType {other:?} for Kraken Spot, skipping");
815                return Ok(());
816            }
817        }
818
819        if let Some(d) = depth {
820            let d_val = d.get();
821            if !matches!(d_val, 10 | 25 | 100 | 500 | 1000) {
822                log::warn!("Invalid depth {d_val} for Kraken Spot, valid: 10, 25, 100, 500, 1000");
823                return Ok(());
824            }
825        }
826
827        let ws = self.ws.clone();
828        self.spawn_ws(
829            async move {
830                ws.subscribe_book(instrument_id, depth.map(|d| d.get() as u32))
831                    .await
832                    .map_err(|e| anyhow::anyhow!("{e}"))
833            },
834            "subscribe book",
835        );
836
837        Ok(())
838    }
839
840    fn subscribe_quotes(&mut self, cmd: SubscribeQuotes) -> anyhow::Result<()> {
841        let instrument_id = cmd.instrument_id;
842        let ws = self.ws.clone();
843
844        self.spawn_ws(
845            async move {
846                ws.subscribe_quotes(instrument_id)
847                    .await
848                    .map_err(|e| anyhow::anyhow!("{e}"))
849            },
850            "subscribe quotes",
851        );
852
853        Ok(())
854    }
855
856    fn subscribe_trades(&mut self, cmd: SubscribeTrades) -> anyhow::Result<()> {
857        let instrument_id = cmd.instrument_id;
858        let ws = self.ws.clone();
859
860        self.spawn_ws(
861            async move {
862                ws.subscribe_trades(instrument_id)
863                    .await
864                    .map_err(|e| anyhow::anyhow!("{e}"))
865            },
866            "subscribe trades",
867        );
868
869        Ok(())
870    }
871
872    fn subscribe_mark_prices(&mut self, cmd: SubscribeMarkPrices) -> anyhow::Result<()> {
873        log::warn!(
874            "Mark price subscription not supported for Spot instrument {}",
875            cmd.instrument_id
876        );
877        Ok(())
878    }
879
880    fn subscribe_index_prices(&mut self, cmd: SubscribeIndexPrices) -> anyhow::Result<()> {
881        log::warn!(
882            "Index price subscription not supported for Spot instrument {}",
883            cmd.instrument_id
884        );
885        Ok(())
886    }
887
888    fn subscribe_bars(&mut self, cmd: SubscribeBars) -> anyhow::Result<()> {
889        let bar_type = cmd.bar_type;
890
891        if bar_type.aggregation_source() != AggregationSource::External {
892            log::warn!("Cannot subscribe to {bar_type} bars: only EXTERNAL bars supported");
893            return Ok(());
894        }
895
896        if !bar_type.spec().is_time_aggregated() {
897            log::warn!("Cannot subscribe to {bar_type} bars: only time-based bars supported");
898            return Ok(());
899        }
900
901        let ws = self.ws.clone();
902        self.spawn_ws(
903            async move {
904                ws.subscribe_bars(bar_type)
905                    .await
906                    .map_err(|e| anyhow::anyhow!("{e}"))
907            },
908            "subscribe bars",
909        );
910
911        Ok(())
912    }
913
914    fn subscribe_instrument_status(
915        &mut self,
916        cmd: SubscribeInstrumentStatus,
917    ) -> anyhow::Result<()> {
918        log::debug!(
919            "subscribe_instrument_status: {} (status changes detected via periodic instrument polling)",
920            cmd.instrument_id,
921        );
922        Ok(())
923    }
924
925    fn unsubscribe_book_deltas(&mut self, cmd: &UnsubscribeBookDeltas) -> anyhow::Result<()> {
926        let instrument_id = cmd.instrument_id;
927
928        if self.ws_l3.as_ref().is_some_and(|ws| {
929            ws.subscriptions_contains(&format!("level3:{}", instrument_id.symbol))
930        }) {
931            let symbol_ustr = instrument_id.symbol.inner();
932
933            if let Some(ws_l3) = self.ws_l3.clone() {
934                self.spawn_ws(
935                    async move {
936                        ws_l3
937                            .unsubscribe_book_l3(symbol_ustr)
938                            .await
939                            .map_err(|e| anyhow::anyhow!("{e}"))?;
940                        Ok(())
941                    },
942                    "unsubscribe l3 book",
943                );
944            }
945            return Ok(());
946        }
947
948        let ws = self.ws.clone();
949        self.spawn_ws(
950            async move {
951                ws.unsubscribe_book(instrument_id)
952                    .await
953                    .map_err(|e| anyhow::anyhow!("{e}"))
954            },
955            "unsubscribe book",
956        );
957
958        Ok(())
959    }
960
961    fn unsubscribe_quotes(&mut self, cmd: &UnsubscribeQuotes) -> anyhow::Result<()> {
962        let instrument_id = cmd.instrument_id;
963        let ws = self.ws.clone();
964
965        self.spawn_ws(
966            async move {
967                ws.unsubscribe_quotes(instrument_id)
968                    .await
969                    .map_err(|e| anyhow::anyhow!("{e}"))
970            },
971            "unsubscribe quotes",
972        );
973
974        Ok(())
975    }
976
977    fn unsubscribe_trades(&mut self, cmd: &UnsubscribeTrades) -> anyhow::Result<()> {
978        let instrument_id = cmd.instrument_id;
979        let ws = self.ws.clone();
980
981        self.spawn_ws(
982            async move {
983                ws.unsubscribe_trades(instrument_id)
984                    .await
985                    .map_err(|e| anyhow::anyhow!("{e}"))
986            },
987            "unsubscribe trades",
988        );
989
990        Ok(())
991    }
992
993    fn unsubscribe_mark_prices(&mut self, _cmd: &UnsubscribeMarkPrices) -> anyhow::Result<()> {
994        Ok(())
995    }
996
997    fn unsubscribe_index_prices(&mut self, _cmd: &UnsubscribeIndexPrices) -> anyhow::Result<()> {
998        Ok(())
999    }
1000
1001    fn unsubscribe_bars(&mut self, cmd: &UnsubscribeBars) -> anyhow::Result<()> {
1002        let bar_type = cmd.bar_type;
1003        let ws = self.ws.clone();
1004
1005        self.spawn_ws(
1006            async move {
1007                ws.unsubscribe_bars(bar_type)
1008                    .await
1009                    .map_err(|e| anyhow::anyhow!("{e}"))
1010            },
1011            "unsubscribe bars",
1012        );
1013
1014        Ok(())
1015    }
1016
1017    fn unsubscribe_instrument_status(
1018        &mut self,
1019        _cmd: &UnsubscribeInstrumentStatus,
1020    ) -> anyhow::Result<()> {
1021        Ok(())
1022    }
1023
1024    fn request_instruments(&self, request: RequestInstruments) -> anyhow::Result<()> {
1025        let http = self.http.clone();
1026        let sender = self.data_sender.clone();
1027        let instruments_cache = self.instruments.clone();
1028        let request_id = request.request_id;
1029        let client_id = request.client_id.unwrap_or(self.client_id);
1030        let venue = *KRAKEN_VENUE;
1031        let start_nanos = datetime_to_unix_nanos(request.start);
1032        let end_nanos = datetime_to_unix_nanos(request.end);
1033        let params = request.params;
1034        let clock = self.clock;
1035
1036        self.spawn_command(async move {
1037            match http.request_instruments(None).await {
1038                Ok(instruments) => {
1039                    instruments_cache.rcu(|m| {
1040                        for instrument in &instruments {
1041                            m.insert(instrument.id(), instrument.clone());
1042                        }
1043                    });
1044                    http.cache_instruments(&instruments);
1045
1046                    let response = DataResponse::Instruments(InstrumentsResponse::new(
1047                        request_id,
1048                        client_id,
1049                        venue,
1050                        instruments,
1051                        start_nanos,
1052                        end_nanos,
1053                        clock.get_time_ns(),
1054                        params,
1055                    ));
1056
1057                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1058                        log::error!("Failed to send instruments response: {e}");
1059                    }
1060                }
1061                Err(e) => log::error!("Instruments request failed: {e:?}"),
1062            }
1063        });
1064
1065        Ok(())
1066    }
1067
1068    fn request_instrument(&self, request: RequestInstrument) -> anyhow::Result<()> {
1069        let http = self.http.clone();
1070        let sender = self.data_sender.clone();
1071        let instruments = self.instruments.clone();
1072        let instrument_id = request.instrument_id;
1073        let request_id = request.request_id;
1074        let client_id = request.client_id.unwrap_or(self.client_id);
1075        let start_nanos = datetime_to_unix_nanos(request.start);
1076        let end_nanos = datetime_to_unix_nanos(request.end);
1077        let params = request.params;
1078        let clock = self.clock;
1079
1080        self.spawn_command(async move {
1081            match http.request_instruments(None).await {
1082                Ok(all_instruments) => {
1083                    instruments.rcu(|m| {
1084                        for instrument in &all_instruments {
1085                            m.insert(instrument.id(), instrument.clone());
1086                        }
1087                    });
1088                    http.cache_instruments(&all_instruments);
1089
1090                    let instrument = all_instruments
1091                        .into_iter()
1092                        .find(|i| i.id() == instrument_id);
1093
1094                    if let Some(instrument) = instrument {
1095                        let response = DataResponse::Instrument(Box::new(InstrumentResponse::new(
1096                            request_id,
1097                            client_id,
1098                            instrument.id(),
1099                            instrument,
1100                            start_nanos,
1101                            end_nanos,
1102                            clock.get_time_ns(),
1103                            params,
1104                        )));
1105
1106                        if let Err(e) = sender.send(DataEvent::Response(response)) {
1107                            log::error!("Failed to send instrument response: {e}");
1108                        }
1109                    } else {
1110                        log::error!("Instrument not found: {instrument_id}");
1111                    }
1112                }
1113                Err(e) => log::error!("Instrument request failed: {e:?}"),
1114            }
1115        });
1116
1117        Ok(())
1118    }
1119    fn request_trades(&self, request: RequestTrades) -> anyhow::Result<()> {
1120        let http = self.http.clone();
1121        let sender = self.data_sender.clone();
1122        let instrument_id = request.instrument_id;
1123        let start = request.start;
1124        let end = request.end;
1125        let limit = request.limit.map(|n| n.get() as u64);
1126        let request_id = request.request_id;
1127        let client_id = request.client_id.unwrap_or(self.client_id);
1128        let params = request.params;
1129        let clock = self.clock;
1130        let start_nanos = datetime_to_unix_nanos(start);
1131        let end_nanos = datetime_to_unix_nanos(end);
1132
1133        self.spawn_command(async move {
1134            match http.request_trades(instrument_id, start, end, limit).await {
1135                Ok(trades) => {
1136                    let response = DataResponse::Trades(TradesResponse::new(
1137                        request_id,
1138                        client_id,
1139                        instrument_id,
1140                        trades,
1141                        start_nanos,
1142                        end_nanos,
1143                        clock.get_time_ns(),
1144                        params,
1145                    ));
1146
1147                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1148                        log::error!("Failed to send trades response: {e}");
1149                    }
1150                }
1151                Err(e) => log::error!("Trades request failed: {e:?}"),
1152            }
1153        });
1154
1155        Ok(())
1156    }
1157
1158    fn request_bars(&self, request: RequestBars) -> anyhow::Result<()> {
1159        let http = self.http.clone();
1160        let sender = self.data_sender.clone();
1161        let bar_type = request.bar_type;
1162        let start = request.start;
1163        let end = request.end;
1164        let limit = request.limit.map(|n| n.get() as u64);
1165        let request_id = request.request_id;
1166        let client_id = request.client_id.unwrap_or(self.client_id);
1167        let params = request.params;
1168        let clock = self.clock;
1169        let start_nanos = datetime_to_unix_nanos(start);
1170        let end_nanos = datetime_to_unix_nanos(end);
1171
1172        self.spawn_command(async move {
1173            match http.request_bars(bar_type, start, end, limit).await {
1174                Ok(bars) => {
1175                    let response = DataResponse::Bars(BarsResponse::new(
1176                        request_id,
1177                        client_id,
1178                        bar_type,
1179                        bars,
1180                        start_nanos,
1181                        end_nanos,
1182                        clock.get_time_ns(),
1183                        params,
1184                    ));
1185
1186                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1187                        log::error!("Failed to send bars response: {e}");
1188                    }
1189                }
1190                Err(e) => log::error!("Bars request failed: {e:?}"),
1191            }
1192        });
1193
1194        Ok(())
1195    }
1196
1197    fn request_book_snapshot(&self, request: RequestBookSnapshot) -> anyhow::Result<()> {
1198        let http = self.http.clone();
1199        let sender = self.data_sender.clone();
1200        let instrument_id = request.instrument_id;
1201        let depth = request.depth.map(|n| n.get() as u32);
1202        let request_id = request.request_id;
1203        let client_id = request.client_id.unwrap_or(self.client_id);
1204        let params = request.params;
1205        let clock = self.clock;
1206
1207        self.spawn_command(async move {
1208            match http.request_book_snapshot(instrument_id, depth).await {
1209                Ok(book) => {
1210                    let response = DataResponse::Book(BookResponse::new(
1211                        request_id,
1212                        client_id,
1213                        instrument_id,
1214                        book,
1215                        None,
1216                        None,
1217                        clock.get_time_ns(),
1218                        params,
1219                    ));
1220
1221                    if let Err(e) = sender.send(DataEvent::Response(response)) {
1222                        log::error!("Failed to send book snapshot response: {e}");
1223                    }
1224                }
1225                Err(e) => log::error!("Book snapshot request failed: {e:?}"),
1226            }
1227        });
1228
1229        Ok(())
1230    }
1231}
1232
1233#[cfg(test)]
1234mod tests {
1235    use nautilus_common::{live::runner::set_data_event_sender, messages::DataEvent};
1236    use nautilus_model::{
1237        enums::{BookAction, RecordFlag},
1238        identifiers::Symbol,
1239        instruments::{InstrumentAny, currency_pair::CurrencyPair},
1240        types::{Currency, Price, Quantity},
1241    };
1242    use rstest::rstest;
1243    use rust_decimal::Decimal;
1244    use rust_decimal_macros::dec;
1245
1246    use super::*;
1247    use crate::{
1248        common::consts::KRAKEN_CLIENT_ID,
1249        config::KrakenDataClientConfig,
1250        websocket::spot_v2::{
1251            level_3::messages::KrakenL3Snapshot,
1252            messages::{KrakenWsBookData, KrakenWsBookLevel},
1253        },
1254    };
1255
1256    fn setup_test_env() {
1257        let (sender, _receiver) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1258        set_data_event_sender(sender);
1259    }
1260
1261    fn make_instrument() -> InstrumentAny {
1262        InstrumentAny::CurrencyPair(
1263            CurrencyPair::builder()
1264                .instrument_id(InstrumentId::from("BTC/USD.KRAKEN"))
1265                .raw_symbol(Symbol::from("BTC/USD"))
1266                .base_currency(Currency::BTC())
1267                .quote_currency(Currency::USD())
1268                .price_precision(1)
1269                .size_precision(8)
1270                .price_increment(Price::from("0.1"))
1271                .size_increment(Quantity::from("0.00000001"))
1272                .ts_event(UnixNanos::default())
1273                .ts_init(UnixNanos::default())
1274                .build()
1275                .unwrap(),
1276        )
1277    }
1278
1279    #[rstest]
1280    fn test_spot_data_client_new() {
1281        setup_test_env();
1282        let config = KrakenDataClientConfig::default();
1283        let client = KrakenSpotDataClient::new(*KRAKEN_CLIENT_ID, config);
1284        assert!(client.is_ok());
1285
1286        let client = client.unwrap();
1287        assert_eq!(client.client_id(), *KRAKEN_CLIENT_ID);
1288        assert_eq!(client.venue(), Some(*KRAKEN_VENUE));
1289        assert!(!client.is_connected());
1290        assert!(client.is_disconnected());
1291        assert!(client.instruments().is_empty());
1292    }
1293
1294    #[rstest]
1295    fn test_l3_snapshot_checksum_mismatch_emits_clear_and_requests_resync() {
1296        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1297        let instruments = Arc::new(AtomicMap::new());
1298        let instrument = make_instrument();
1299        instruments.insert(instrument.id(), instrument);
1300
1301        let depths = Arc::new(Mutex::new(AHashMap::new()));
1302        depths.lock().insert("BTC/USD".to_string(), 1000);
1303
1304        let snapshot: KrakenL3Snapshot = serde_json::from_str(
1305            r#"{
1306                "symbol": "BTC/USD",
1307                "bids": [{
1308                    "order_id": "order-bid-1",
1309                    "limit_price": 4199.0,
1310                    "order_qty": 3.00000000,
1311                    "timestamp": "2024-01-01T00:00:00Z"
1312                }],
1313                "asks": [{
1314                    "order_id": "order-ask-1",
1315                    "limit_price": 4200.0,
1316                    "order_qty": 0.01000000,
1317                    "timestamp": "2024-01-01T00:00:00Z"
1318                }],
1319                "checksum": 1,
1320                "timestamp": "2024-01-01T00:00:00Z"
1321            }"#,
1322        )
1323        .unwrap();
1324
1325        let mut states = AHashMap::new();
1326        let hasher = BookOrderIdHasher::new();
1327        let mut sink = DataEventSink { sender: &sender };
1328        let request = process_l3_message(
1329            KrakenL3WsMessage::Snapshot(snapshot),
1330            &mut sink,
1331            &instruments,
1332            &depths,
1333            &mut states,
1334            &hasher,
1335            true,
1336            get_atomic_clock_realtime().get_time_ns(),
1337        )
1338        .expect("expected resync request");
1339
1340        assert_eq!(request.symbol, "BTC/USD");
1341        assert_eq!(request.depth, 1000);
1342        assert_eq!(request.reason, "snapshot checksum mismatch");
1343
1344        let event = receiver.try_recv().expect("expected clear event");
1345        let DataEvent::Data(Data::Deltas(deltas)) = event else {
1346            panic!("expected deltas event");
1347        };
1348
1349        assert_eq!(deltas.deltas.len(), 1);
1350        assert_eq!(deltas.deltas[0].action, BookAction::Clear);
1351        assert!(states["BTC/USD"].awaiting_snapshot);
1352        assert!(states["BTC/USD"].open_orders.is_empty());
1353        assert!(receiver.try_recv().is_err());
1354    }
1355
1356    #[rstest]
1357    fn test_l2_update_prunes_levels_beyond_subscribed_depth() {
1358        let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel::<DataEvent>();
1359        let instruments = Arc::new(AtomicMap::new());
1360        let instrument = make_instrument();
1361        let instrument_id = instrument.id();
1362        instruments.insert(instrument_id, instrument);
1363
1364        let book_sequence = Arc::new(AtomicU64::new(0));
1365        let l2_depths = L2Depths::default();
1366        l2_depths.insert("BTC/USD", 10);
1367        let mut l2_books = L2BookState::default();
1368        let ohlc_buffer = Arc::new(Mutex::new(AHashMap::new()));
1369        let context = SpotMessageContext {
1370            sender: &sender,
1371            instruments: &instruments,
1372            book_sequence: &book_sequence,
1373            l2_depths: &l2_depths,
1374            ohlc_buffer: &ohlc_buffer,
1375            clock: get_atomic_clock_realtime(),
1376        };
1377
1378        let snapshot = KrakenWsBookData {
1379            symbol: Ustr::from("BTC/USD"),
1380            bids: Some(
1381                (0..10)
1382                    .map(|i| book_level(Decimal::from(100 - i), Decimal::ONE))
1383                    .collect(),
1384            ),
1385            asks: Some(
1386                (0..10)
1387                    .map(|i| book_level(Decimal::from(101 + i), Decimal::ONE))
1388                    .collect(),
1389            ),
1390            checksum: Some(0),
1391            timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
1392        };
1393        KrakenSpotDataClient::handle_ws_message(
1394            KrakenSpotWsMessage::Book {
1395                data: vec![snapshot],
1396                is_snapshot: true,
1397            },
1398            &context,
1399            &mut l2_books,
1400        );
1401
1402        let DataEvent::Data(Data::Deltas(snapshot_deltas)) =
1403            receiver.try_recv().expect("expected snapshot deltas")
1404        else {
1405            panic!("expected snapshot deltas");
1406        };
1407        assert_eq!(snapshot_deltas.deltas.len(), 21);
1408        assert_eq!(snapshot_deltas.deltas[0].action, BookAction::Clear);
1409        assert!(RecordFlag::F_LAST.matches(snapshot_deltas.deltas.last().unwrap().flags));
1410
1411        let bid_update = KrakenWsBookData {
1412            symbol: Ustr::from("BTC/USD"),
1413            bids: Some(vec![book_level(dec!(100.5), Decimal::ONE)]),
1414            asks: Some(vec![]),
1415            checksum: Some(0),
1416            timestamp: "2024-01-01T00:00:01Z".parse().unwrap(),
1417        };
1418        KrakenSpotDataClient::handle_ws_message(
1419            KrakenSpotWsMessage::Book {
1420                data: vec![bid_update],
1421                is_snapshot: false,
1422            },
1423            &context,
1424            &mut l2_books,
1425        );
1426
1427        let DataEvent::Data(Data::Deltas(bid_update_deltas)) =
1428            receiver.try_recv().expect("expected bid update deltas")
1429        else {
1430            panic!("expected bid update deltas");
1431        };
1432        assert_eq!(bid_update_deltas.deltas.len(), 2);
1433        assert_eq!(bid_update_deltas.deltas[0].action, BookAction::Update);
1434        assert_eq!(bid_update_deltas.deltas[1].action, BookAction::Delete);
1435        assert_eq!(bid_update_deltas.deltas[1].order.price, Price::from("91.0"));
1436        assert!(RecordFlag::F_LAST.matches(bid_update_deltas.deltas[1].flags));
1437
1438        let ask_update = KrakenWsBookData {
1439            symbol: Ustr::from("BTC/USD"),
1440            bids: Some(vec![]),
1441            asks: Some(vec![book_level(dec!(100.6), Decimal::ONE)]),
1442            checksum: Some(0),
1443            timestamp: "2024-01-01T00:00:02Z".parse().unwrap(),
1444        };
1445        KrakenSpotDataClient::handle_ws_message(
1446            KrakenSpotWsMessage::Book {
1447                data: vec![ask_update],
1448                is_snapshot: false,
1449            },
1450            &context,
1451            &mut l2_books,
1452        );
1453
1454        let DataEvent::Data(Data::Deltas(ask_update_deltas)) =
1455            receiver.try_recv().expect("expected ask update deltas")
1456        else {
1457            panic!("expected ask update deltas");
1458        };
1459        assert_eq!(ask_update_deltas.deltas.len(), 2);
1460        assert_eq!(ask_update_deltas.deltas[0].action, BookAction::Update);
1461        assert_eq!(ask_update_deltas.deltas[1].action, BookAction::Delete);
1462        assert_eq!(
1463            ask_update_deltas.deltas[1].order.price,
1464            Price::from("110.0")
1465        );
1466        assert!(RecordFlag::F_LAST.matches(ask_update_deltas.deltas[1].flags));
1467
1468        let book = l2_books
1469            .books
1470            .get(&instrument_id)
1471            .expect("expected shadow book");
1472        assert_eq!(book.bids(None).count(), 10);
1473        assert_eq!(book.asks(None).count(), 10);
1474        assert_eq!(book.best_bid_price(), Some(Price::from("100.5")));
1475        assert_eq!(book.best_ask_price(), Some(Price::from("100.6")));
1476        assert!(receiver.try_recv().is_err());
1477    }
1478
1479    #[rstest]
1480    fn test_spot_data_client_start_stop() {
1481        setup_test_env();
1482        let config = KrakenDataClientConfig::default();
1483        let mut client = KrakenSpotDataClient::new(*KRAKEN_CLIENT_ID, config).unwrap();
1484
1485        assert!(client.start().is_ok());
1486        assert!(client.stop().is_ok());
1487        assert!(client.is_disconnected());
1488    }
1489
1490    fn book_level(price: Decimal, qty: Decimal) -> KrakenWsBookLevel {
1491        KrakenWsBookLevel { price, qty }
1492    }
1493}