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