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