Skip to main content

nautilus_kraken/websocket/spot_v2/
parse.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! WebSocket message parsers for converting Kraken streaming data to Nautilus domain models.
17
18use 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
43/// Parses Kraken WebSocket ticker data into a Nautilus quote tick.
44///
45/// # Errors
46///
47/// Returns an error if:
48/// - Bid or ask price/quantity cannot be parsed.
49pub 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
87/// Parses Kraken WebSocket trade data into a Nautilus trade tick.
88///
89/// # Errors
90///
91/// Returns an error if:
92/// - Price or quantity cannot be parsed.
93/// - Timestamp is invalid.
94pub 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
128/// Parses Kraken WebSocket book data into Nautilus order book deltas.
129///
130/// Returns a vector of deltas, one for each bid and ask level.
131///
132/// # Errors
133///
134/// Returns an error if:
135/// - Price or quantity cannot be parsed.
136/// - Timestamp is invalid.
137pub 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
288/// Parses Kraken WebSocket OHLC data into a Nautilus bar.
289///
290/// The bar's `ts_event` is computed as `interval_begin` + `interval` minutes.
291///
292/// # Errors
293///
294/// Returns an error if:
295/// - Price or quantity values cannot be parsed.
296/// - The interval cannot be converted to a valid bar specification.
297pub 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    // Compute bar close time: interval_begin + interval minutes
316    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
323/// Converts a Kraken OHLC interval (minutes) to a Nautilus bar specification.
324fn 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), // 21600 min = 360 hours = 15 days
335        _ => anyhow::bail!("Unsupported Kraken OHLC interval: {interval}"),
336    };
337
338    Ok(BarSpecification::new(step, aggregation, PriceType::Last))
339}
340
341/// Parses Kraken execution type and order status to Nautilus order status.
342fn 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
366/// Parses Kraken order type to Nautilus order type.
367fn 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        // Trailing stops lack offset fields in WS reports, map to non-trailing equivalents
376        Some(KrakenOrderType::TrailingStop) => OrderType::StopMarket,
377        Some(KrakenOrderType::TrailingStopLimit) => OrderType::StopLimit,
378        Some(KrakenOrderType::SettlePosition) => OrderType::Market,
379        None => OrderType::Limit,
380    }
381}
382
383/// Parses Kraken time-in-force to Nautilus time-in-force.
384fn parse_time_in_force(
385    time_in_force: Option<KrakenTimeInForce>,
386    post_only: Option<bool>,
387) -> TimeInForce {
388    // Handle post_only flag
389    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
406/// Parses a Kraken WebSocket execution message into an [`OrderStatusReport`].
407///
408/// # Errors
409///
410/// Returns an error if required fields are missing or cannot be parsed.
411pub 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    // Quantity fallback: order_qty -> cached -> cum_qty -> last_qty (for trade snapshots)
429    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, // client_order_id set below if present
457        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    // Price fallback: limit_price -> avg_price -> last_price
477    // Note: pending_new messages may not include any price fields, which is fine for
478    // orders we submitted (engine already has the price from submission)
479    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    // avg_px fallback: avg_price -> cum_cost / cum_qty -> last_price (for single trades/snapshots)
492    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    // Set trigger type for conditional orders (WebSocket doesn't provide trigger field)
520    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
535/// Parses a Kraken WebSocket trade execution into a [`FillReport`].
536///
537/// This should only be called when exec_type is "trade".
538///
539/// # Errors
540///
541/// Returns an error if required fields are missing or cannot be parsed.
542pub 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    // Calculate commission from fees array
583    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, // venue_position_id
614        ts_event,
615        ts_init,
616        None, // report_id
617    ))
618}
619
620/// Parses a raw WebSocket JSON string and returns [`KrakenSpotWsMessage::OrderResponse`] if the
621/// message is an order-method response envelope, or `Ok(None)` for unrecognized messages.
622///
623/// # Errors
624///
625/// Returns an error if the message appears to be an order response but cannot be deserialized.
626pub 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        // Verify ts_event is computed as interval_begin + interval (close time)
992        // interval_begin is 2023-10-04T16:25:00Z, interval is 1 minute, so close is 16:26:00Z
993        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), // 21600 min = 15 days
1011        ];
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}