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, 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#[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 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 #[must_use]
169 pub fn instruments(&self) -> Vec<InstrumentAny> {
170 self.instruments.load().values().cloned().collect()
171 }
172
173 #[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}