1use anyhow::Context;
19use jiff::Timestamp;
20use nautilus_core::{UUID4, nanos::UnixNanos};
21use nautilus_model::{
22 data::{Bar, BarSpecification, BarType, BookOrder, OrderBookDelta, QuoteTick, TradeTick},
23 enums::{
24 AggregationSource, AggressorSide, BarAggregation, BookAction, LiquiditySide, OrderSide,
25 OrderStatus, OrderType, PriceType, RecordFlag, TimeInForce, TriggerType,
26 },
27 identifiers::{AccountId, ClientOrderId, InstrumentId, TradeId, VenueOrderId},
28 instruments::{Instrument, any::InstrumentAny},
29 reports::{FillReport, OrderStatusReport},
30 types::{Currency, Money, Price, Quantity},
31};
32use rust_decimal::Decimal;
33
34use super::{
35 enums::{KrakenExecType, KrakenLiquidityInd, KrakenWsOrderStatus},
36 messages::{
37 KrakenSpotWsMessage, KrakenWsBookData, KrakenWsBookLevel, KrakenWsExecutionData,
38 KrakenWsOhlcData, KrakenWsOrderResponse, KrakenWsTickerData, KrakenWsTradeData,
39 },
40};
41use crate::common::enums::{KrakenOrderSide, KrakenOrderType, KrakenTimeInForce};
42
43pub fn parse_quote_tick(
50 ticker: &KrakenWsTickerData,
51 instrument: &InstrumentAny,
52 ts_init: UnixNanos,
53) -> anyhow::Result<QuoteTick> {
54 let instrument_id = instrument.id();
55 let price_precision = instrument.price_precision();
56 let size_precision = instrument.size_precision();
57
58 let bid_price = Price::from_decimal_dp(ticker.bid, price_precision).with_context(|| {
59 format!("Failed to construct bid Price with precision {price_precision}")
60 })?;
61 let bid_size =
62 Quantity::from_decimal_dp(ticker.bid_qty, size_precision).with_context(|| {
63 format!("Failed to construct bid Quantity with precision {size_precision}")
64 })?;
65
66 let ask_price = Price::from_decimal_dp(ticker.ask, price_precision).with_context(|| {
67 format!("Failed to construct ask Price with precision {price_precision}")
68 })?;
69 let ask_size =
70 Quantity::from_decimal_dp(ticker.ask_qty, size_precision).with_context(|| {
71 format!("Failed to construct ask Quantity with precision {size_precision}")
72 })?;
73
74 let ts_event = datetime_to_nanos(ticker.timestamp, "ticker.timestamp")?;
75
76 Ok(QuoteTick::new(
77 instrument_id,
78 bid_price,
79 ask_price,
80 bid_size,
81 ask_size,
82 ts_event,
83 ts_init,
84 ))
85}
86
87pub fn parse_trade_tick(
95 trade: &KrakenWsTradeData,
96 instrument: &InstrumentAny,
97 ts_init: UnixNanos,
98) -> anyhow::Result<TradeTick> {
99 let instrument_id = instrument.id();
100 let price_precision = instrument.price_precision();
101 let size_precision = instrument.size_precision();
102
103 let price = Price::from_decimal_dp(trade.price, price_precision)
104 .with_context(|| format!("Failed to construct Price with precision {price_precision}"))?;
105 let size = Quantity::from_decimal_dp(trade.qty, size_precision)
106 .with_context(|| format!("Failed to construct Quantity with precision {size_precision}"))?;
107
108 let aggressor = match trade.side {
109 KrakenOrderSide::Buy => AggressorSide::Buy,
110 KrakenOrderSide::Sell => AggressorSide::Sell,
111 };
112
113 let trade_id = TradeId::new_checked(trade.trade_id.to_string())?;
114 let ts_event = datetime_to_nanos(trade.timestamp, "trade.timestamp")?;
115
116 TradeTick::new_checked(
117 instrument_id,
118 price,
119 size,
120 aggressor,
121 trade_id,
122 ts_event,
123 ts_init,
124 )
125 .context("Failed to construct TradeTick from Kraken WebSocket trade")
126}
127
128pub fn parse_book_deltas(
138 book: &KrakenWsBookData,
139 instrument: &InstrumentAny,
140 sequence: u64,
141 is_snapshot: bool,
142 ts_init: UnixNanos,
143) -> anyhow::Result<Vec<OrderBookDelta>> {
144 let instrument_id = instrument.id();
145 let price_precision = instrument.price_precision();
146 let size_precision = instrument.size_precision();
147
148 let ts_event = datetime_to_nanos(book.timestamp, "book.timestamp")?;
149
150 let mut current_sequence = sequence;
151 let mut deltas = Vec::new();
152
153 if is_snapshot {
154 deltas.push(OrderBookDelta::clear(
155 instrument_id,
156 current_sequence,
157 ts_event,
158 ts_init,
159 ));
160 current_sequence += 1;
161 }
162
163 if let Some(ref bids) = book.bids {
164 parse_book_side_levels(
165 bids,
166 OrderSide::Buy,
167 is_snapshot,
168 instrument_id,
169 price_precision,
170 size_precision,
171 ts_event,
172 ts_init,
173 &mut current_sequence,
174 &mut deltas,
175 )?;
176 }
177
178 if let Some(ref asks) = book.asks {
179 parse_book_side_levels(
180 asks,
181 OrderSide::Sell,
182 is_snapshot,
183 instrument_id,
184 price_precision,
185 size_precision,
186 ts_event,
187 ts_init,
188 &mut current_sequence,
189 &mut deltas,
190 )?;
191 }
192
193 if let Some(last) = deltas.last_mut() {
194 last.flags |= RecordFlag::F_LAST as u8;
195 }
196
197 Ok(deltas)
198}
199
200#[expect(clippy::too_many_arguments)]
201fn parse_book_side_levels(
202 levels: &[KrakenWsBookLevel],
203 side: OrderSide,
204 is_snapshot: bool,
205 instrument_id: InstrumentId,
206 price_precision: u8,
207 size_precision: u8,
208 ts_event: UnixNanos,
209 ts_init: UnixNanos,
210 current_sequence: &mut u64,
211 deltas: &mut Vec<OrderBookDelta>,
212) -> anyhow::Result<()> {
213 for level in levels {
214 let Some(delta) = parse_book_level(
215 level,
216 side,
217 is_snapshot,
218 instrument_id,
219 price_precision,
220 size_precision,
221 *current_sequence,
222 ts_event,
223 ts_init,
224 )?
225 else {
226 continue;
227 };
228 deltas.push(delta);
229 *current_sequence += 1;
230 }
231
232 Ok(())
233}
234
235#[expect(clippy::too_many_arguments)]
236fn parse_book_level(
237 level: &KrakenWsBookLevel,
238 side: OrderSide,
239 is_snapshot: bool,
240 instrument_id: InstrumentId,
241 price_precision: u8,
242 size_precision: u8,
243 sequence: u64,
244 ts_event: UnixNanos,
245 ts_init: UnixNanos,
246) -> anyhow::Result<Option<OrderBookDelta>> {
247 let price = Price::from_decimal_dp(level.price, price_precision)
248 .with_context(|| format!("Failed to construct Price with precision {price_precision}"))?;
249 let size = Quantity::from_decimal_dp(level.qty, size_precision)
250 .with_context(|| format!("Failed to construct Quantity with precision {size_precision}"))?;
251
252 let action = if is_snapshot {
253 if size.raw == 0 {
254 return Ok(None);
255 }
256 BookAction::Add
257 } else if size.raw == 0 {
258 BookAction::Delete
259 } else {
260 BookAction::Update
261 };
262
263 let order_id = price.raw as u64;
264 let order = BookOrder::new(side, price, size, order_id);
265 let mut flags = RecordFlag::F_MBP as u8;
266 if is_snapshot {
267 flags |= RecordFlag::F_SNAPSHOT as u8;
268 }
269
270 Ok(Some(OrderBookDelta::new(
271 instrument_id,
272 action,
273 order,
274 flags,
275 sequence,
276 ts_event,
277 ts_init,
278 )))
279}
280
281pub(super) fn datetime_to_nanos(value: Timestamp, field: &str) -> anyhow::Result<UnixNanos> {
282 let nanos = u64::try_from(value.as_nanosecond())
283 .with_context(|| format!("Timestamp predates Unix epoch: {field}='{value}'"))?;
284 Ok(UnixNanos::from(nanos))
285}
286
287pub fn parse_ws_bar(
297 ohlc: &KrakenWsOhlcData,
298 instrument: &InstrumentAny,
299 ts_init: UnixNanos,
300) -> anyhow::Result<Bar> {
301 let instrument_id = instrument.id();
302 let price_precision = instrument.price_precision();
303 let size_precision = instrument.size_precision();
304
305 let open = Price::from_decimal_dp(ohlc.open, price_precision)?;
306 let high = Price::from_decimal_dp(ohlc.high, price_precision)?;
307 let low = Price::from_decimal_dp(ohlc.low, price_precision)?;
308 let close = Price::from_decimal_dp(ohlc.close, price_precision)?;
309 let volume = Quantity::from_decimal_dp(ohlc.volume, size_precision)?;
310
311 let bar_spec = interval_to_bar_spec(ohlc.interval)?;
312 let bar_type = BarType::new(instrument_id, bar_spec, AggregationSource::External);
313
314 let interval_secs = i64::from(ohlc.interval) * 60;
316 let close_time = ohlc.interval_begin + jiff::SignedDuration::from_secs(interval_secs);
317 let ts_event = UnixNanos::from(u64::try_from(close_time.as_nanosecond()).unwrap_or(0));
318
319 Bar::new_checked(bar_type, open, high, low, close, volume, ts_event, ts_init)
320}
321
322fn interval_to_bar_spec(interval: u32) -> anyhow::Result<BarSpecification> {
324 let (step, aggregation) = match interval {
325 1 => (1, BarAggregation::Minute),
326 5 => (5, BarAggregation::Minute),
327 15 => (15, BarAggregation::Minute),
328 30 => (30, BarAggregation::Minute),
329 60 => (1, BarAggregation::Hour),
330 240 => (4, BarAggregation::Hour),
331 1440 => (1, BarAggregation::Day),
332 10080 => (1, BarAggregation::Week),
333 21600 => (15, BarAggregation::Day), _ => anyhow::bail!("Unsupported Kraken OHLC interval: {interval}"),
335 };
336
337 Ok(BarSpecification::new(step, aggregation, PriceType::Last))
338}
339
340fn parse_order_status(
342 exec_type: KrakenExecType,
343 order_status: Option<KrakenWsOrderStatus>,
344) -> OrderStatus {
345 match exec_type {
346 KrakenExecType::Canceled => return OrderStatus::Canceled,
347 KrakenExecType::Expired => return OrderStatus::Expired,
348 KrakenExecType::Filled => return OrderStatus::Filled,
349 KrakenExecType::Trade => {
350 return match order_status {
351 Some(KrakenWsOrderStatus::Filled) => OrderStatus::Filled,
352 Some(KrakenWsOrderStatus::PartiallyFilled) | None => OrderStatus::PartiallyFilled,
353 Some(status) => status.into(),
354 };
355 }
356 _ => {}
357 }
358
359 match order_status {
360 Some(status) => status.into(),
361 None => OrderStatus::Accepted,
362 }
363}
364
365fn parse_order_type(order_type: Option<KrakenOrderType>) -> OrderType {
367 match order_type {
368 Some(KrakenOrderType::Market) => OrderType::Market,
369 Some(KrakenOrderType::Limit) => OrderType::Limit,
370 Some(KrakenOrderType::StopLoss) => OrderType::StopMarket,
371 Some(KrakenOrderType::TakeProfit) => OrderType::MarketIfTouched,
372 Some(KrakenOrderType::StopLossLimit) => OrderType::StopLimit,
373 Some(KrakenOrderType::TakeProfitLimit) => OrderType::LimitIfTouched,
374 Some(KrakenOrderType::TrailingStop) => OrderType::StopMarket,
376 Some(KrakenOrderType::TrailingStopLimit) => OrderType::StopLimit,
377 Some(KrakenOrderType::SettlePosition) => OrderType::Market,
378 None => OrderType::Limit,
379 }
380}
381
382fn parse_time_in_force(
384 time_in_force: Option<KrakenTimeInForce>,
385 post_only: Option<bool>,
386) -> TimeInForce {
387 if post_only == Some(true) {
389 return TimeInForce::Gtc;
390 }
391
392 match time_in_force {
393 Some(KrakenTimeInForce::GoodTilCancelled) => TimeInForce::Gtc,
394 Some(KrakenTimeInForce::ImmediateOrCancel) => TimeInForce::Ioc,
395 Some(KrakenTimeInForce::GoodTilDate) => TimeInForce::Gtd,
396 Some(KrakenTimeInForce::FillOrKill) => TimeInForce::Fok,
397 None => TimeInForce::Gtc,
398 }
399}
400
401fn parse_liquidity_side(liquidity_ind: Option<KrakenLiquidityInd>) -> LiquiditySide {
402 liquidity_ind.map_or(LiquiditySide::NoLiquiditySide, Into::into)
403}
404
405pub fn parse_ws_order_status_report(
411 exec: &KrakenWsExecutionData,
412 instrument: &InstrumentAny,
413 account_id: AccountId,
414 cached_order_qty: Option<Decimal>,
415 ts_init: UnixNanos,
416) -> anyhow::Result<OrderStatusReport> {
417 let instrument_id = instrument.id();
418 let venue_order_id = VenueOrderId::new(&exec.order_id);
419 let order_side = exec.side.map(Into::into);
420 let order_type = parse_order_type(exec.order_type);
421 let time_in_force = parse_time_in_force(exec.time_in_force, exec.post_only);
422 let order_status = parse_order_status(exec.exec_type, exec.order_status);
423
424 let price_precision = instrument.price_precision();
425 let size_precision = instrument.size_precision();
426
427 let last_qty = exec
429 .last_qty
430 .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
431 .transpose()
432 .context("Failed to parse last_qty")?;
433
434 let filled_qty = exec
435 .cum_qty
436 .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
437 .transpose()
438 .context("Failed to parse cum_qty")?
439 .or(last_qty)
440 .unwrap_or_else(|| Quantity::zero(size_precision));
441
442 let quantity = exec
443 .order_qty
444 .or(cached_order_qty)
445 .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
446 .transpose()
447 .context("Failed to parse order_qty")?
448 .unwrap_or(filled_qty);
449
450 let ts_event = datetime_to_nanos(exec.timestamp, "execution.timestamp")?;
451
452 let mut report = OrderStatusReport::new(
453 account_id,
454 instrument_id,
455 None, venue_order_id,
457 order_side,
458 order_type,
459 time_in_force,
460 order_status,
461 quantity,
462 filled_qty,
463 ts_event,
464 ts_event,
465 ts_init,
466 Some(UUID4::new()),
467 );
468
469 if let Some(ref cl_ord_id) = exec.cl_ord_id
470 && !cl_ord_id.is_empty()
471 {
472 report = report.with_client_order_id(ClientOrderId::new(cl_ord_id));
473 }
474
475 let price_value = exec
479 .limit_price
480 .filter(|p| *p > Decimal::ZERO)
481 .or(exec.avg_price.filter(|p| *p > Decimal::ZERO))
482 .or(exec.last_price.filter(|p| *p > Decimal::ZERO));
483
484 if let Some(px) = price_value {
485 let price =
486 Price::from_decimal_dp(px, price_precision).context("Failed to parse order price")?;
487 report = report.with_price(price);
488 }
489
490 let avg_px = exec
492 .avg_price
493 .filter(|p| *p > Decimal::ZERO)
494 .or_else(|| match (exec.cum_cost, exec.cum_qty) {
495 (Some(cost), Some(qty)) if qty > Decimal::ZERO => Some(cost / qty),
496 _ => None,
497 })
498 .or_else(|| exec.last_price.filter(|p| *p > Decimal::ZERO));
499
500 if let Some(avg_price) = avg_px {
501 report.avg_px = Some(avg_price);
502 }
503
504 if exec.post_only == Some(true) {
505 report = report.with_post_only(true);
506 }
507
508 if exec.reduce_only == Some(true) {
509 report = report.with_reduce_only(true);
510 }
511
512 if let Some(ref reason) = exec.reason
513 && !reason.is_empty()
514 {
515 report = report.with_cancel_reason(reason.clone());
516 }
517
518 let is_conditional = matches!(
520 order_type,
521 OrderType::StopMarket
522 | OrderType::StopLimit
523 | OrderType::MarketIfTouched
524 | OrderType::LimitIfTouched
525 );
526
527 if is_conditional {
528 report = report.with_trigger_type(TriggerType::Default);
529 }
530
531 Ok(report)
532}
533
534pub fn parse_ws_fill_report(
542 exec: &KrakenWsExecutionData,
543 instrument: &InstrumentAny,
544 account_id: AccountId,
545 ts_init: UnixNanos,
546) -> anyhow::Result<FillReport> {
547 let instrument_id = instrument.id();
548 let venue_order_id = VenueOrderId::new(&exec.order_id);
549
550 let exec_id = exec
551 .exec_id
552 .as_ref()
553 .context("Missing exec_id for trade execution")?;
554 let trade_id =
555 TradeId::new_checked(exec_id).context("Invalid exec_id in Kraken trade execution")?;
556
557 let order_side = exec
558 .side
559 .map(Into::into)
560 .context("Missing side for trade execution")?;
561
562 let price_precision = instrument.price_precision();
563 let size_precision = instrument.size_precision();
564
565 let last_qty = exec
566 .last_qty
567 .map(|qty| Quantity::from_decimal_dp(qty, size_precision))
568 .transpose()
569 .context("Failed to parse last_qty")?
570 .context("Missing last_qty for trade execution")?;
571
572 let last_px = exec
573 .last_price
574 .map(|px| Price::from_decimal_dp(px, price_precision))
575 .transpose()
576 .context("Failed to parse last_price")?
577 .context("Missing last_price for trade execution")?;
578
579 let liquidity_side = parse_liquidity_side(exec.liquidity_ind);
580
581 let commission = if let Some(ref fees) = exec.fees {
583 if let Some(fee) = fees.first() {
584 let currency = Currency::get_or_create_crypto(&fee.asset);
585 Money::from_decimal(fee.qty.abs(), currency).context("Failed to parse fill fee")?
586 } else {
587 Money::zero(instrument.quote_currency())
588 }
589 } else {
590 Money::zero(instrument.quote_currency())
591 };
592
593 let ts_event = datetime_to_nanos(exec.timestamp, "execution.timestamp")?;
594
595 let client_order_id = exec
596 .cl_ord_id
597 .as_ref()
598 .filter(|s| !s.is_empty())
599 .map(ClientOrderId::new);
600
601 Ok(FillReport::new(
602 account_id,
603 instrument_id,
604 venue_order_id,
605 trade_id,
606 order_side,
607 last_qty,
608 last_px,
609 commission,
610 liquidity_side,
611 client_order_id,
612 None, ts_event,
614 ts_init,
615 None, ))
617}
618
619pub fn parse_order_response(text: &str) -> anyhow::Result<Option<KrakenSpotWsMessage>> {
626 let value: serde_json::Value =
627 serde_json::from_str(text).with_context(|| format!("Failed to parse JSON: {text}"))?;
628
629 let method_str = match value.get("method").and_then(|m| m.as_str()) {
630 Some(s) => s.to_owned(),
631 None => return Ok(None),
632 };
633
634 if !matches!(
635 method_str.as_str(),
636 "add_order" | "amend_order" | "cancel_order" | "batch_add"
637 ) {
638 return Ok(None);
639 }
640
641 let response: KrakenWsOrderResponse = serde_json::from_value(value).with_context(|| {
642 format!("Failed to deserialize order response for method '{method_str}'")
643 })?;
644 Ok(Some(KrakenSpotWsMessage::OrderResponse(response)))
645}
646
647#[cfg(test)]
648mod tests {
649 use nautilus_model::{identifiers::Symbol, types::Currency};
650 use rstest::rstest;
651 use rust_decimal_macros::dec;
652 use ustr::Ustr;
653
654 use super::*;
655 use crate::{common::consts::KRAKEN_VENUE, websocket::spot_v2::messages::KrakenWsRawMessage};
656
657 const TS: UnixNanos = UnixNanos::new(1_700_000_000_000_000_000);
658
659 #[rstest]
660 fn test_parse_time_in_force_fok() {
661 assert_eq!(
662 parse_time_in_force(Some(KrakenTimeInForce::FillOrKill), None),
663 TimeInForce::Fok
664 );
665 }
666
667 fn load_test_json(filename: &str) -> String {
668 let path = format!("test_data/{filename}");
669 std::fs::read_to_string(&path)
670 .unwrap_or_else(|e| panic!("Failed to load test data from {path}: {e}"))
671 }
672
673 fn create_mock_instrument() -> InstrumentAny {
674 use nautilus_model::instruments::currency_pair::CurrencyPair;
675
676 let instrument_id = InstrumentId::new(Symbol::new("BTC/USD"), *KRAKEN_VENUE);
677 InstrumentAny::CurrencyPair(
678 CurrencyPair::builder()
679 .instrument_id(instrument_id)
680 .raw_symbol(Symbol::new("XBTUSDT"))
681 .base_currency(Currency::BTC())
682 .quote_currency(Currency::USDT())
683 .price_precision(1)
684 .size_precision(8)
685 .price_increment(Price::from("0.1"))
686 .size_increment(Quantity::from("0.00000001"))
687 .ts_event(TS)
688 .ts_init(TS)
689 .build()
690 .unwrap(),
691 )
692 }
693
694 #[rstest]
695 fn test_parse_quote_tick() {
696 let json = load_test_json("ws_ticker_snapshot.json");
697 let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
698 let ticker: KrakenWsTickerData = serde_json::from_str(message.data[0].get()).unwrap();
699
700 let instrument = create_mock_instrument();
701 let quote_tick = parse_quote_tick(&ticker, &instrument, TS).unwrap();
702
703 assert_eq!(quote_tick.instrument_id, instrument.id());
704 assert_eq!(quote_tick.bid_price, Price::from("105944.20"));
705 assert_eq!(quote_tick.ask_price, Price::from("105944.30"));
706 assert_eq!(quote_tick.bid_size, Quantity::from("2.5"));
707 assert_eq!(quote_tick.ask_size, Quantity::from("3.2"));
708 assert_eq!(
709 quote_tick.ts_event,
710 UnixNanos::from(1_671_960_659_123_456_000)
711 );
712 assert_eq!(quote_tick.ts_init, TS);
713 }
714
715 #[rstest]
716 fn test_parse_trade_tick() {
717 let json = load_test_json("ws_trade_update.json");
718 let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
719 let trade: KrakenWsTradeData = serde_json::from_str(message.data[0].get()).unwrap();
720
721 let instrument = create_mock_instrument();
722 let trade_tick = parse_trade_tick(&trade, &instrument, TS).unwrap();
723
724 assert_eq!(trade_tick.instrument_id, instrument.id());
725 assert_eq!(trade_tick.price, Price::from("105944.20"));
726 assert_eq!(trade_tick.size, Quantity::from("0.00027625"));
727 assert!(matches!(
728 trade_tick.aggressor_side,
729 AggressorSide::Buy | AggressorSide::Sell
730 ));
731 assert_eq!(
732 trade_tick.ts_event,
733 UnixNanos::from(1_696_613_755_440_295_000)
734 );
735 assert_eq!(trade_tick.ts_init, TS);
736 }
737
738 #[rstest]
739 fn test_parse_book_deltas_snapshot() {
740 let json = load_test_json("ws_book_snapshot.json");
741 let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
742 let book: KrakenWsBookData = serde_json::from_str(message.data[0].get()).unwrap();
743
744 let instrument = create_mock_instrument();
745 let deltas = parse_book_deltas(&book, &instrument, 1, true, TS).unwrap();
746
747 assert!(!deltas.is_empty());
748
749 let bid_count = deltas
750 .iter()
751 .filter(|d| d.order.side == OrderSide::Buy.into())
752 .count();
753 let ask_count = deltas
754 .iter()
755 .filter(|d| d.order.side == OrderSide::Sell.into())
756 .count();
757
758 assert!(bid_count > 0);
759 assert!(ask_count > 0);
760
761 let first_delta = &deltas[0];
762 assert_eq!(first_delta.instrument_id, instrument.id());
763 assert_eq!(first_delta.action, BookAction::Clear);
764 assert!(RecordFlag::F_SNAPSHOT.matches(first_delta.flags));
765 assert!(!RecordFlag::F_LAST.matches(first_delta.flags));
766
767 assert!(deltas[1..].iter().all(|d| d.action == BookAction::Add));
768 assert!(
769 deltas[1..]
770 .iter()
771 .all(|d| RecordFlag::F_MBP.matches(d.flags))
772 );
773 assert!(
774 deltas[1..]
775 .iter()
776 .all(|d| RecordFlag::F_SNAPSHOT.matches(d.flags))
777 );
778 assert!(RecordFlag::F_LAST.matches(deltas.last().unwrap().flags));
779
780 let expected_ts_event = UnixNanos::from(1_696_613_755_440_295_000);
781 assert!(deltas.iter().all(|d| d.ts_event == expected_ts_event));
782 assert!(deltas.iter().all(|d| d.ts_init == TS));
783 }
784
785 #[rstest]
786 fn test_parse_book_deltas_update() {
787 let json = load_test_json("ws_book_update.json");
788 let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
789 let book: KrakenWsBookData = serde_json::from_str(message.data[0].get()).unwrap();
790
791 let instrument = create_mock_instrument();
792 let deltas = parse_book_deltas(&book, &instrument, 1, false, TS).unwrap();
793
794 assert!(!deltas.is_empty());
795
796 let first_delta = &deltas[0];
797 assert_eq!(first_delta.instrument_id, instrument.id());
798 assert_eq!(first_delta.action, BookAction::Update);
799 assert_eq!(first_delta.order.side, OrderSide::Buy.into());
800 assert_eq!(first_delta.order.price, Price::from("105944.20"));
801 assert!(RecordFlag::F_MBP.matches(first_delta.flags));
802 assert!(RecordFlag::F_LAST.matches(first_delta.flags));
803 assert!(!RecordFlag::F_SNAPSHOT.matches(first_delta.flags));
804
805 let expected_ts_event = UnixNanos::from(1_696_613_755_440_295_000);
806 assert!(deltas.iter().all(|d| d.ts_event == expected_ts_event));
807 assert!(deltas.iter().all(|d| d.ts_init == TS));
808 }
809
810 #[rstest]
811 fn test_parse_book_deltas_snapshot_skips_zero_qty_levels() {
812 let book = KrakenWsBookData {
813 symbol: Ustr::from("BTC/USD"),
814 bids: Some(vec![KrakenWsBookLevel {
815 price: dec!(100),
816 qty: Decimal::ZERO,
817 }]),
818 asks: Some(vec![KrakenWsBookLevel {
819 price: dec!(101),
820 qty: dec!(2),
821 }]),
822 checksum: Some(0),
823 timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
824 };
825
826 let instrument = create_mock_instrument();
827 let deltas = parse_book_deltas(&book, &instrument, 7, true, TS).unwrap();
828
829 assert_eq!(deltas.len(), 2);
830 assert_eq!(deltas[0].action, BookAction::Clear);
831 assert_eq!(deltas[0].sequence, 7);
832 assert!(RecordFlag::F_SNAPSHOT.matches(deltas[0].flags));
833 assert!(!RecordFlag::F_LAST.matches(deltas[0].flags));
834
835 let add = &deltas[1];
836 assert_eq!(add.action, BookAction::Add);
837 assert_eq!(add.sequence, 8);
838 assert_eq!(add.order.side, OrderSide::Sell.into());
839 assert_eq!(add.order.price, Price::from("101.0"));
840 assert!(RecordFlag::F_MBP.matches(add.flags));
841 assert!(RecordFlag::F_SNAPSHOT.matches(add.flags));
842 assert!(RecordFlag::F_LAST.matches(add.flags));
843 }
844
845 #[rstest]
846 fn test_parse_book_deltas_update_zero_qty_deletes_level() {
847 let book = KrakenWsBookData {
848 symbol: Ustr::from("BTC/USD"),
849 bids: Some(vec![KrakenWsBookLevel {
850 price: dec!(100),
851 qty: Decimal::ZERO,
852 }]),
853 asks: Some(vec![]),
854 checksum: Some(0),
855 timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
856 };
857
858 let instrument = create_mock_instrument();
859 let deltas = parse_book_deltas(&book, &instrument, 11, false, TS).unwrap();
860
861 assert_eq!(deltas.len(), 1);
862 let delete = &deltas[0];
863 assert_eq!(delete.action, BookAction::Delete);
864 assert_eq!(delete.sequence, 11);
865 assert_eq!(delete.order.side, OrderSide::Buy.into());
866 assert_eq!(delete.order.price, Price::from("100.0"));
867 assert_eq!(delete.order.size.raw, 0);
868 assert!(RecordFlag::F_MBP.matches(delete.flags));
869 assert!(RecordFlag::F_LAST.matches(delete.flags));
870 assert!(!RecordFlag::F_SNAPSHOT.matches(delete.flags));
871 }
872
873 #[rstest]
874 fn test_parse_ws_order_status_report_preserves_decimal_avg_px() {
875 let execution = ws_execution_data(Some(KrakenOrderSide::Buy));
876
877 let report = parse_ws_order_status_report(
878 &execution,
879 &create_mock_instrument(),
880 AccountId::from("KRAKEN-001"),
881 None,
882 TS,
883 )
884 .unwrap();
885
886 assert_eq!(report.avg_px, Some(dec!(0.1234567890123456789012345678)));
887 }
888
889 #[rstest]
890 fn test_parse_ws_order_status_report_preserves_missing_side() {
891 let execution = ws_execution_data(None);
892
893 let report = parse_ws_order_status_report(
894 &execution,
895 &create_mock_instrument(),
896 AccountId::from("KRAKEN-001"),
897 None,
898 TS,
899 )
900 .unwrap();
901
902 assert_eq!(report.order_side, None);
903 }
904
905 #[rstest]
906 fn test_parse_ws_fill_report_rejects_missing_side() {
907 let mut execution = ws_execution_data(None);
908 execution.exec_type = KrakenExecType::Trade;
909 execution.exec_id = Some("TRADE-1".to_string());
910 execution.last_qty = Some(dec!(1));
911 execution.last_price = Some(dec!(100));
912
913 let error = parse_ws_fill_report(
914 &execution,
915 &create_mock_instrument(),
916 AccountId::from("KRAKEN-001"),
917 TS,
918 )
919 .expect_err("a trade execution without a side must be rejected");
920
921 assert_eq!(error.to_string(), "Missing side for trade execution");
922 }
923
924 fn ws_execution_data(side: Option<KrakenOrderSide>) -> KrakenWsExecutionData {
925 KrakenWsExecutionData {
926 exec_type: KrakenExecType::Status,
927 order_id: "ORDER-1".to_string(),
928 cl_ord_id: Some("CLIENT-1".to_string()),
929 symbol: Some("BTC/USD".to_string()),
930 side,
931 order_type: Some(KrakenOrderType::Limit),
932 order_qty: Some(dec!(3)),
933 limit_price: None,
934 order_status: Some(KrakenWsOrderStatus::PartiallyFilled),
935 cum_qty: Some(dec!(3)),
936 cum_cost: Some(dec!(0.3703703670370370367037037034)),
937 avg_price: None,
938 time_in_force: Some(KrakenTimeInForce::GoodTilCancelled),
939 post_only: Some(false),
940 reduce_only: Some(false),
941 timestamp: "2024-01-01T00:00:00Z".parse().unwrap(),
942 exec_id: None,
943 last_qty: None,
944 last_price: None,
945 cost: None,
946 liquidity_ind: None,
947 fees: None,
948 fee_usd_equiv: None,
949 reason: None,
950 }
951 }
952
953 #[rstest]
954 fn test_datetime_to_nanos() {
955 let dt = "2023-10-06T17:35:55.440295Z".parse::<Timestamp>().unwrap();
956 let result = datetime_to_nanos(dt, "test").unwrap();
957 assert_eq!(result, UnixNanos::from(1_696_613_755_440_295_000));
958 }
959
960 #[rstest]
961 fn test_datetime_to_nanos_out_of_range_errors() {
962 let dt = "1500-01-01T00:00:00Z".parse::<Timestamp>().unwrap();
963 let result = datetime_to_nanos(dt, "test");
964 assert!(result.is_err());
965 let err = result.unwrap_err().to_string();
966 assert!(err.contains("test"));
967 }
968
969 #[rstest]
970 fn test_parse_ws_bar() {
971 let json = load_test_json("ws_ohlc_update.json");
972 let message: KrakenWsRawMessage = serde_json::from_str(&json).unwrap();
973 let ohlc: KrakenWsOhlcData = serde_json::from_str(message.data[0].get()).unwrap();
974
975 let instrument = create_mock_instrument();
976 let bar = parse_ws_bar(&ohlc, &instrument, TS).unwrap();
977
978 assert_eq!(bar.bar_type.instrument_id(), instrument.id());
979 assert_eq!(bar.open, Price::from("106038.2"));
980 assert_eq!(bar.high, Price::from("106044.3"));
981 assert_eq!(bar.low, Price::from("106038.1"));
982 assert_eq!(bar.close, Price::from("106040.1"));
983 assert_eq!(bar.volume, Quantity::from("30927.68066226"));
984
985 let spec = bar.bar_type.spec();
986 assert_eq!(spec.step.get(), 1);
987 assert_eq!(spec.aggregation, BarAggregation::Minute);
988 assert_eq!(spec.price_type, PriceType::Last);
989
990 let expected_close = ohlc.interval_begin + jiff::SignedDuration::from_mins(1);
993 let expected_ts_event =
994 UnixNanos::from(u64::try_from(expected_close.as_nanosecond()).unwrap());
995 assert_eq!(bar.ts_event, expected_ts_event);
996 }
997
998 #[rstest]
999 fn test_interval_to_bar_spec() {
1000 let test_cases = [
1001 (1, 1, BarAggregation::Minute),
1002 (5, 5, BarAggregation::Minute),
1003 (15, 15, BarAggregation::Minute),
1004 (30, 30, BarAggregation::Minute),
1005 (60, 1, BarAggregation::Hour),
1006 (240, 4, BarAggregation::Hour),
1007 (1440, 1, BarAggregation::Day),
1008 (10080, 1, BarAggregation::Week),
1009 (21600, 15, BarAggregation::Day), ];
1011
1012 for (interval, expected_step, expected_aggregation) in test_cases {
1013 let spec = interval_to_bar_spec(interval).unwrap();
1014 assert_eq!(
1015 spec.step.get(),
1016 expected_step,
1017 "Failed for interval {interval}"
1018 );
1019 assert_eq!(
1020 spec.aggregation, expected_aggregation,
1021 "Failed for interval {interval}"
1022 );
1023 assert_eq!(spec.price_type, PriceType::Last);
1024 }
1025 }
1026
1027 #[rstest]
1028 fn test_interval_to_bar_spec_invalid() {
1029 let result = interval_to_bar_spec(999);
1030 assert!(result.is_err());
1031 }
1032
1033 #[rstest]
1034 fn test_parse_order_response_envelope_returns_order_response_variant() {
1035 use crate::websocket::spot_v2::enums::KrakenWsMethod;
1036
1037 let raw = load_test_json("ws_add_order_response_success.json");
1038 let parsed = parse_order_response(&raw).expect("parse ok");
1039 match parsed {
1040 Some(KrakenSpotWsMessage::OrderResponse(resp)) => {
1041 assert_eq!(resp.method, KrakenWsMethod::AddOrder);
1042 assert_eq!(resp.req_id, Some(42));
1043 assert!(resp.success);
1044 }
1045 other => panic!("expected OrderResponse, was {other:?}"),
1046 }
1047 }
1048
1049 #[rstest]
1050 fn test_parse_order_response_returns_none_for_non_order_method() {
1051 let json = r#"{"method":"subscribe","req_id":1,"success":true}"#;
1052 let result = parse_order_response(json).expect("parse ok");
1053 assert!(result.is_none());
1054 }
1055
1056 #[rstest]
1057 fn test_parse_order_response_returns_none_for_data_message() {
1058 let json = r#"{"channel":"ticker","type":"snapshot","data":[]}"#;
1059 let result = parse_order_response(json).expect("parse ok");
1060 assert!(result.is_none());
1061 }
1062}