1use 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
83struct 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#[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 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 #[must_use]
178 pub fn instruments(&self) -> Vec<InstrumentAny> {
179 self.instruments.load().values().cloned().collect()
180 }
181
182 #[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}